From aa8c2959ca8476f269786e1317fb6d2938abd3f9 Mon Sep 17 00:00:00 2001 From: Tim Baek Date: Mon, 9 Feb 2026 08:07:33 +0400 Subject: [PATCH] refac --- backend/open_webui/utils/middleware.py | 930 +++++++++++-------------- backend/open_webui/utils/misc.py | 112 ++- 2 files changed, 482 insertions(+), 560 deletions(-) diff --git a/backend/open_webui/utils/middleware.py b/backend/open_webui/utils/middleware.py index 402898ad7..58ea3f249 100644 --- a/backend/open_webui/utils/middleware.py +++ b/backend/open_webui/utils/middleware.py @@ -2777,8 +2777,21 @@ async def non_streaming_chat_response_handler(response, ctx): title = Chats.get_chat_title_by_id(metadata["chat_id"]) - # Use output from backend if provided (OR-compliant backends) + # Use output from backend if provided (OR-compliant backends), + # otherwise generate from response content response_output = response_data.get("output") + if not response_output: + response_output = [ + { + "type": "message", + "id": output_id("msg"), + "status": "completed", + "role": "assistant", + "content": [ + {"type": "output_text", "text": content} + ], + } + ] await event_emitter( { @@ -2786,11 +2799,7 @@ async def non_streaming_chat_response_handler(response, ctx): "data": { "done": True, "content": content, - **( - {"output": response_output} - if response_output - else {} - ), + "output": response_output, "title": title, }, } @@ -2803,7 +2812,7 @@ async def non_streaming_chat_response_handler(response, ctx): { "role": "assistant", "content": content, - **({"output": response_output} if response_output else {}), + "output": response_output, }, ) @@ -2878,326 +2887,54 @@ async def streaming_chat_response_handler(response, ctx): # Handle as a background task async def response_handler(response, events): - def serialize_content_blocks(content_blocks, raw=False): - content = "" - - for block in content_blocks: - if block["type"] == "text": - block_content = block["content"].strip() - if block_content: - content = f"{content}{block_content}\n" - elif block["type"] == "tool_calls": - attributes = block.get("attributes", {}) - - tool_calls = block.get("content", []) - results = block.get("results", []) - - if content and not content.endswith("\n"): - content += "\n" - - if results: - - tool_calls_display_content = "" - for tool_call in tool_calls: - - tool_call_id = tool_call.get("id", "") - tool_name = tool_call.get("function", {}).get( - "name", "" - ) - tool_arguments = tool_call.get("function", {}).get( - "arguments", "" - ) - - tool_result = None - tool_result_files = None - for result in results: - if tool_call_id == result.get("tool_call_id", ""): - tool_result = result.get("content", None) - tool_result_files = result.get("files", None) - break - - if tool_result is not None: - tool_result_embeds = result.get("embeds", "") - tool_calls_display_content = f'{tool_calls_display_content}
\nTool Executed\n
\n' - else: - tool_calls_display_content = f'{tool_calls_display_content}
\nExecuting...\n
\n' - - if not raw: - content = f"{content}{tool_calls_display_content}" - else: - tool_calls_display_content = "" - - for tool_call in tool_calls: - tool_call_id = tool_call.get("id", "") - tool_name = tool_call.get("function", {}).get( - "name", "" - ) - tool_arguments = tool_call.get("function", {}).get( - "arguments", "" - ) - - tool_calls_display_content = f'{tool_calls_display_content}\n
\nExecuting...\n
\n' - - if not raw: - content = f"{content}{tool_calls_display_content}" - - elif block["type"] == "reasoning": - reasoning_display_content = html.escape( - "\n".join( - (f"> {line}" if not line.startswith(">") else line) - for line in block["content"].splitlines() - ) - ) - - reasoning_duration = block.get("duration", None) - - start_tag = block.get("start_tag", "") - end_tag = block.get("end_tag", "") - - if content and not content.endswith("\n"): - content += "\n" - - if reasoning_duration is not None: - if raw: - content = ( - f'{content}{start_tag}{block["content"]}{end_tag}\n' - ) - else: - content = f'{content}
\nThought for {reasoning_duration} seconds\n{reasoning_display_content}\n
\n' - else: - if raw: - content = ( - f'{content}{start_tag}{block["content"]}{end_tag}\n' - ) - else: - content = f'{content}
\nThinking…\n{reasoning_display_content}\n
\n' - - elif block["type"] == "code_interpreter": - attributes = block.get("attributes", {}) - output = block.get("output", None) - lang = attributes.get("lang", "") - - content_stripped, original_whitespace = ( - split_content_and_whitespace(content) - ) - if is_opening_code_block(content_stripped): - # Remove trailing backticks that would open a new block - content = ( - content_stripped.rstrip("`").rstrip() - + original_whitespace - ) - else: - # Keep content as is - either closing backticks or no backticks - content = content_stripped + original_whitespace - - if content and not content.endswith("\n"): - content += "\n" - - if output: - output = html.escape(json.dumps(output)) - - if raw: - content = f'{content}\n{block["content"]}\n\n```output\n{output}\n```\n' - else: - content = f'{content}
\nAnalyzed\n```{lang}\n{block["content"]}\n```\n
\n' - else: - if raw: - content = f'{content}\n{block["content"]}\n\n' - else: - content = f'{content}
\nAnalyzing...\n```{lang}\n{block["content"]}\n```\n
\n' - - else: - block_content = str(block["content"]).strip() - if block_content: - content = f"{content}{block['type']}: {block_content}\n" - - return content.strip() - - return content.strip() - - def convert_content_blocks_to_messages(content_blocks, raw=False): - messages = [] - - temp_blocks = [] - for idx, block in enumerate(content_blocks): - if block["type"] == "tool_calls": - messages.append( - { - "role": "assistant", - "content": serialize_content_blocks(temp_blocks, raw), - "tool_calls": block.get("content"), - } - ) - - results = block.get("results", []) - - for result in results: - messages.append( - { - "role": "tool", - "tool_call_id": result["tool_call_id"], - "content": result.get("content", "") or "", - } - ) - temp_blocks = [] - else: - temp_blocks.append(block) - - if temp_blocks: - content = serialize_content_blocks(temp_blocks, raw) - if content: - messages.append( - { - "role": "assistant", - "content": content, - } - ) - - return messages - - def convert_content_blocks_to_output(content_blocks): + def tag_output_handler(content_type, tags, content, output): """ - Convert content_blocks to Open Responses-aligned output items. - See: https://openresponses.org/specification + 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. """ - output_items = [] - - def next_id(prefix): - return f"{prefix}_{uuid4().hex[:24]}" - - for block in content_blocks: - block_type = block.get("type", "") - # Use backend-provided ID if available, fallback to generated - block_id = block.get("id") - - if block_type == "text": - text_content = block.get("content", "").strip() - if text_content: - output_items.append( - { - "type": "message", - "id": block_id or next_id("msg"), - "status": "completed", - "role": "assistant", - "content": [ - {"type": "output_text", "text": text_content} - ], - } - ) - - elif block_type == "tool_calls": - tool_calls = block.get("content", []) - results = block.get("results", []) - - # Emit function_call items - for tool_call in tool_calls: - call_id = tool_call.get("id", "") - func = tool_call.get("function", {}) - output_items.append( - { - "type": "function_call", - "id": call_id - or next_id( - "fc" - ), # Use call_id as item id if available - "call_id": call_id, - "name": func.get("name", ""), - "arguments": func.get("arguments", "{}"), - "status": "completed" if results else "in_progress", - } - ) - - # Emit function_call_output items - for result in results: - output_items.append( - { - "type": "function_call_output", - "id": result.get("id") or next_id("fco"), - "call_id": result.get("tool_call_id", ""), - "output": [ - { - "type": "input_text", - "text": result.get("content", ""), - } - ], - "status": "completed", - **( - {"files": result.get("files")} - if result.get("files") - else {} - ), - **( - {"embeds": result.get("embeds")} - if result.get("embeds") - else {} - ), - } - ) - - elif block_type == "reasoning": - reasoning_content = block.get("content", "").strip() - duration = block.get("duration") - output_items.append( - { - "type": "reasoning", - "id": block_id or next_id("r"), - "status": ( - "completed" - if duration is not None - else "in_progress" - ), - "content": ( - [{"type": "output_text", "text": reasoning_content}] - if reasoning_content - else None - ), - "summary": None, - } - ) - - elif block_type == "code_interpreter": - code = block.get("content", "") - output_val = block.get("output") - attrs = block.get("attributes", {}) - output_items.append( - { - "type": "open_webui:code_interpreter", - "id": block_id or next_id("ci"), - "status": ( - "completed" - if output_val is not None - else "in_progress" - ), - "lang": attrs.get("lang", ""), - "code": code, - "output": output_val, - } - ) - - return output_items - - def tag_content_handler(content_type, tags, content, content_blocks): end_flag = False def extract_attributes(tag_content): """Extract attributes from a tag if they exist.""" attributes = {} - if not tag_content: # Ensure tag_content is not None + if not tag_content: return attributes - # Match attributes in the format: key="value" (ignores single quotes for simplicity) matches = re.findall(r'(\w+)\s*=\s*"([^"]+)"', tag_content) for key, value in matches: attributes[key] = value return attributes - if content_blocks[-1]["type"] == "text": + def get_last_text(out): + """Get text from last message item, or empty string.""" + if out and out[-1].get("type") == "message": + parts = out[-1].get("content", []) + if parts and parts[-1].get("type") == "output_text": + return parts[-1].get("text", "") + return "" + + def set_last_text(out, text): + """Set text on last message item's output_text.""" + if out and out[-1].get("type") == "message": + parts = out[-1].get("content", []) + if parts and parts[-1].get("type") == "output_text": + parts[-1]["text"] = text + + # Map content_type to output item type + output_type_map = { + "reasoning": "reasoning", + "solution": "message", # solution tags just produce text + "code_interpreter": "open_webui:code_interpreter", + } + output_item_type = output_type_map.get(content_type, content_type) + + last_type = output[-1].get("type", "") if output else "" + + if last_type == "message": for start_tag, end_tag in tags: start_tag_pattern = rf"{re.escape(start_tag)}" if start_tag.startswith("<") and start_tag.endswith(">"): - # Match start tag e.g., or - # remove both '<' and '>' from start_tag - # Match start tag with attributes start_tag_pattern = ( rf"<{re.escape(start_tag[1:-1])}(\s.*?)?>" ) @@ -3207,70 +2944,128 @@ async def streaming_chat_response_handler(response, ctx): try: attr_content = ( match.group(1) if match.group(1) else "" - ) # Ensure it's not None + ) except: attr_content = "" - attributes = extract_attributes( - attr_content - ) # Extract attributes safely + attributes = extract_attributes(attr_content) - # Capture everything before and after the matched tag - before_tag = content[ - : match.start() - ] # Content before opening tag - after_tag = content[ - match.end() : - ] # Content after opening tag + before_tag = content[: match.start()] + after_tag = content[match.end() :] - # Remove the start tag and after from the currently handling text block - content_blocks[-1]["content"] = content_blocks[-1][ - "content" - ].replace(match.group(0) + after_tag, "") - - if before_tag: - content_blocks[-1]["content"] = before_tag - - if not content_blocks[-1]["content"]: - content_blocks.pop() - - # Append the new block - content_blocks.append( - { - "type": content_type, - "start_tag": start_tag, - "end_tag": end_tag, - "attributes": attributes, - "content": "", - "started_at": time.time(), - } + # 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, "") ) + if before_tag: + set_last_text(output, before_tag) + + if not get_last_text(output).strip(): + # Remove empty message item + if output and output[-1].get("type") == "message": + output.pop() + + # Append the new output item + if output_item_type == "reasoning": + output.append( + { + "type": "reasoning", + "id": output_id("r"), + "status": "in_progress", + "start_tag": start_tag, + "end_tag": end_tag, + "attributes": attributes, + "content": [], + "summary": None, + "started_at": time.time(), + } + ) + elif output_item_type == "open_webui:code_interpreter": + output.append( + { + "type": "open_webui:code_interpreter", + "id": output_id("ci"), + "status": "in_progress", + "start_tag": start_tag, + "end_tag": end_tag, + "attributes": attributes, + "lang": attributes.get("lang", "python"), + "code": "", + "output": None, + "started_at": time.time(), + } + ) + else: + # solution or other text-producing tag + output.append( + { + "type": "message", + "id": output_id("msg"), + "status": "in_progress", + "role": "assistant", + "content": [{"type": "output_text", "text": ""}], + "_tag_type": content_type, + "start_tag": start_tag, + "end_tag": end_tag, + "attributes": attributes, + "started_at": time.time(), + } + ) + if after_tag: - content_blocks[-1]["content"] = after_tag - tag_content_handler( - content_type, tags, after_tag, content_blocks + # Set the after_tag content on the new item + if output_item_type == "reasoning": + output[-1]["content"] = [ + {"type": "output_text", "text": after_tag} + ] + elif output_item_type == "open_webui:code_interpreter": + output[-1]["code"] = after_tag + else: + set_last_text(output, after_tag) + + tag_output_handler( + content_type, tags, after_tag, output ) break - elif content_blocks[-1]["type"] == content_type: - start_tag = content_blocks[-1]["start_tag"] - end_tag = content_blocks[-1]["end_tag"] + + elif ( + (last_type == "reasoning" and content_type == "reasoning") + or (last_type == "open_webui:code_interpreter" and content_type == "code_interpreter") + or (last_type == "message" and output[-1].get("_tag_type") == content_type) + ): + item = output[-1] + start_tag = item.get("start_tag", "") + end_tag = item.get("end_tag", "") if end_tag.startswith("<") and end_tag.endswith(">"): - # Match end tag e.g., end_tag_pattern = rf"{re.escape(end_tag)}" else: - # Handle cases where end_tag is just a tag name end_tag_pattern = rf"{re.escape(end_tag)}" - # Check if the content has the end tag if re.search(end_tag_pattern, content): end_flag = True - block_content = content_blocks[-1]["content"] - # Strip start and end tags from the content - start_tag_pattern = rf"<{re.escape(start_tag)}(.*?)>" + # 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)}" + if start_tag.startswith("<") and start_tag.endswith(">"): + start_tag_pattern = ( + rf"<{re.escape(start_tag[1:-1])}(\s.*?)?>" + ) block_content = re.sub( start_tag_pattern, "", block_content ).strip() @@ -3278,79 +3073,98 @@ async def streaming_chat_response_handler(response, ctx): end_tag_regex = re.compile(end_tag_pattern, re.DOTALL) split_content = end_tag_regex.split(block_content, maxsplit=1) - # Content inside the tag block_content = ( split_content[0].strip() if split_content else "" ) - - # Leftover content (everything after ``) leftover_content = ( split_content[1].strip() if len(split_content) > 1 else "" ) if block_content: - content_blocks[-1]["content"] = block_content - content_blocks[-1]["ended_at"] = time.time() - content_blocks[-1]["duration"] = int( - content_blocks[-1]["ended_at"] - - content_blocks[-1]["started_at"] - ) + # Update the item with final content + if last_type == "reasoning": + item["content"] = [ + {"type": "output_text", "text": block_content} + ] + item["ended_at"] = time.time() + item["duration"] = int( + item["ended_at"] - item["started_at"] + ) + item["status"] = "completed" + elif last_type == "open_webui:code_interpreter": + item["code"] = block_content + item["ended_at"] = time.time() + item["duration"] = int( + item["ended_at"] - item["started_at"] + ) + else: + set_last_text(output, block_content) + item["ended_at"] = time.time() - # Reset the content_blocks by appending a new text block + # Reset by appending a new message item for leftover if content_type != "code_interpreter": - if leftover_content: - - content_blocks.append( - { - "type": "text", - "content": leftover_content, - } - ) - else: - content_blocks.append( - { - "type": "text", - "content": "", - } - ) - - else: - # Remove the block if content is empty - content_blocks.pop() - - if leftover_content: - content_blocks.append( + output.append( { - "type": "text", - "content": leftover_content, + "type": "message", + "id": output_id("msg"), + "status": "in_progress", + "role": "assistant", + "content": [ + { + "type": "output_text", + "text": leftover_content, + } + ], } ) else: - content_blocks.append( + output.append( { - "type": "text", - "content": "", + "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() + output.append( + { + "type": "message", + "id": output_id("msg"), + "status": "in_progress", + "role": "assistant", + "content": [ + { + "type": "output_text", + "text": leftover_content, + } + ], + } + ) # Clean processed content - start_tag_pattern = rf"{re.escape(start_tag)}" + start_tag_clean = rf"{re.escape(start_tag)}" if start_tag.startswith("<") and start_tag.endswith(">"): - # Match start tag e.g., or - # remove both '<' and '>' from start_tag - # Match start tag with attributes - start_tag_pattern = ( + start_tag_clean = ( rf"<{re.escape(start_tag[1:-1])}(\s.*?)?>" ) content = re.sub( - rf"{start_tag_pattern}(.|\n)*?{re.escape(end_tag)}", + rf"{start_tag_clean}(.|\n)*?{re.escape(end_tag)}", "", content, flags=re.DOTALL, ) - return content, content_blocks, end_flag + return content, output, end_flag message = Chats.get_message_by_id_and_message_id( metadata["chat_id"], metadata["message_id"] @@ -3392,13 +3206,7 @@ async def streaming_chat_response_handler(response, ctx): else: output = [] - # Keep content_blocks for backward compatibility during transition - content_blocks = [ - { - "type": "text", - "content": content, - } - ] + usage = None reasoning_tags_param = metadata.get("params", {}).get("reasoning_tags") @@ -3439,7 +3247,6 @@ async def streaming_chat_response_handler(response, ctx): async def stream_body_handler(response, form_data): nonlocal content - nonlocal content_blocks nonlocal usage nonlocal output @@ -3677,19 +3484,26 @@ async def streaming_chat_response_handler(response, ctx): # Flush any pending text first await flush_pending_delta_data() - pending_content_blocks = content_blocks + [ - { - "type": "tool_calls", - "content": response_tool_calls, - "pending": True, - } - ] + # Build pending function_call output items for display + pending_fc_items = [] + for tc in response_tool_calls: + call_id = tc.get("id", "") + func = tc.get("function", {}) + pending_fc_items.append({ + "type": "function_call", + "id": call_id or output_id("fc"), + "call_id": call_id, + "name": func.get("name", ""), + "arguments": func.get("arguments", "{}"), + "status": "in_progress", + }) + pending_output = output + pending_fc_items await event_emitter( { "type": "chat:completion", "data": { - "content": serialize_content_blocks( - pending_content_blocks + "content": serialize_output( + pending_output ), }, } @@ -3724,52 +3538,64 @@ async def streaming_chat_response_handler(response, ctx): ) if reasoning_content: if ( - not content_blocks - or content_blocks[-1]["type"] != "reasoning" + not output + or output[-1].get("type") != "reasoning" ): - reasoning_block = { + reasoning_item = { "type": "reasoning", + "id": output_id("r"), + "status": "in_progress", "start_tag": "", "end_tag": "", "attributes": { "type": "reasoning_content" }, - "content": "", + "content": [], + "summary": None, "started_at": time.time(), } - content_blocks.append(reasoning_block) + output.append(reasoning_item) else: - reasoning_block = content_blocks[-1] + reasoning_item = output[-1] - reasoning_block["content"] += reasoning_content + # Append to reasoning content + parts = reasoning_item.get("content", []) + if parts and parts[-1].get("type") == "output_text": + parts[-1]["text"] += reasoning_content + else: + reasoning_item["content"] = [ + {"type": "output_text", "text": reasoning_content} + ] data = { - "content": serialize_content_blocks( - content_blocks - ) + "content": serialize_output(output) } if value: if ( - content_blocks - and content_blocks[-1]["type"] + output + and output[-1].get("type") == "reasoning" - and content_blocks[-1] + and output[-1] .get("attributes", {}) .get("type") == "reasoning_content" ): - reasoning_block = content_blocks[-1] - reasoning_block["ended_at"] = time.time() - reasoning_block["duration"] = int( - reasoning_block["ended_at"] - - reasoning_block["started_at"] + reasoning_item = output[-1] + reasoning_item["ended_at"] = time.time() + reasoning_item["duration"] = int( + reasoning_item["ended_at"] + - reasoning_item["started_at"] ) + reasoning_item["status"] = "completed" - content_blocks.append( + output.append( { - "type": "text", - "content": "", + "type": "message", + "id": output_id("msg"), + "status": "in_progress", + "role": "assistant", + "content": [{"type": "output_text", "text": ""}], } ) @@ -3789,44 +3615,55 @@ async def streaming_chat_response_handler(response, ctx): ) content = f"{content}{value}" - if not content_blocks: - content_blocks.append( + if ( + not output + or output[-1].get("type") != "message" + ): + output.append( { - "type": "text", - "content": "", + "type": "message", + "id": output_id("msg"), + "status": "in_progress", + "role": "assistant", + "content": [{"type": "output_text", "text": ""}], } ) - content_blocks[-1]["content"] = ( - content_blocks[-1]["content"] + 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: + output[-1]["content"] = [ + {"type": "output_text", "text": value} + ] if DETECT_REASONING_TAGS: - content, content_blocks, _ = ( - tag_content_handler( + content, output, _ = ( + tag_output_handler( "reasoning", reasoning_tags, content, - content_blocks, + output, ) ) - content, content_blocks, _ = ( - tag_content_handler( + content, output, _ = ( + tag_output_handler( "solution", DEFAULT_SOLUTION_TAGS, content, - content_blocks, + output, ) ) if DETECT_CODE_INTERPRETER: - content, content_blocks, end = ( - tag_content_handler( + content, output, end = ( + tag_output_handler( "code_interpreter", DEFAULT_CODE_INTERPRETER_TAGS, content, - content_blocks, + output, ) ) @@ -3835,9 +3672,6 @@ async def streaming_chat_response_handler(response, ctx): if ENABLE_REALTIME_CHAT_SAVE: # Save message in the database - output = convert_content_blocks_to_output( - content_blocks - ) Chats.upsert_message_to_chat_by_id_and_message_id( metadata["chat_id"], metadata["message_id"], @@ -3848,8 +3682,8 @@ async def streaming_chat_response_handler(response, ctx): ) else: data = { - "content": serialize_content_blocks( - content_blocks + "content": serialize_output( + output ), } @@ -3874,32 +3708,36 @@ async def streaming_chat_response_handler(response, ctx): continue await flush_pending_delta_data() - if content_blocks: - # Clean up the last text block - if content_blocks[-1]["type"] == "text": - content_blocks[-1]["content"] = content_blocks[-1][ - "content" - ].strip() + if output: + # Clean up the last message item + if output[-1].get("type") == "message": + parts = output[-1].get("content", []) + if parts and parts[-1].get("type") == "output_text": + parts[-1]["text"] = parts[-1]["text"].strip() - if not content_blocks[-1]["content"]: - content_blocks.pop() + if not parts[-1]["text"]: + output.pop() - if not content_blocks: - content_blocks.append( - { - "type": "text", - "content": "", - } - ) + if not output: + output.append( + { + "type": "message", + "id": output_id("msg"), + "status": "in_progress", + "role": "assistant", + "content": [{"type": "output_text", "text": ""}], + } + ) - if content_blocks[-1]["type"] == "reasoning": - reasoning_block = content_blocks[-1] - if reasoning_block.get("ended_at") is None: - reasoning_block["ended_at"] = time.time() - reasoning_block["duration"] = int( - reasoning_block["ended_at"] - - reasoning_block["started_at"] + if output[-1].get("type") == "reasoning": + reasoning_item = output[-1] + if reasoning_item.get("ended_at") is None: + reasoning_item["ended_at"] = time.time() + reasoning_item["duration"] = int( + reasoning_item["ended_at"] + - reasoning_item["started_at"] ) + reasoning_item["status"] = "completed" if response_tool_calls: tool_calls.append(response_tool_calls) @@ -3921,14 +3759,19 @@ async def streaming_chat_response_handler(response, ctx): response_tool_calls = tool_calls.pop(0) - content_blocks.append( - { - "type": "tool_calls", - "content": response_tool_calls, - } - ) + # Append function_call items for each tool call + for tc in response_tool_calls: + call_id = tc.get("id", "") + func = tc.get("function", {}) + output.append({ + "type": "function_call", + "id": call_id or output_id("fc"), + "call_id": call_id, + "name": func.get("name", ""), + "arguments": func.get("arguments", "{}"), + "status": "in_progress", + }) - output = convert_content_blocks_to_output(content_blocks) await event_emitter( { "type": "chat:completion", @@ -3964,11 +3807,7 @@ async def streaming_chat_response_handler(response, ctx): f"Error parsing tool call arguments: {tool_args}" ) - # Mutate the original tool call response params as they are passed back to the passed - # back to the LLM via the content blocks. If they are in a json block and are invalid json, - # this can cause downstream LLM integrations to fail (e.g. bedrock gateway) where response - # params are not valid json. - # Main case so far is no args = "" = invalid json. + # Ensure arguments are valid JSON for downstream LLM integrations log.debug( f"Parsed args from {tool_args} to {tool_function_params}" ) @@ -4085,11 +3924,49 @@ async def streaming_chat_response_handler(response, ctx): } ) - content_blocks[-1]["results"] = results - content_blocks.append( + # Update function_call statuses and append function_call_output items + for tc in response_tool_calls: + call_id = tc.get("id", "") + # Mark function_call as completed + for item in output: + if item.get("type") == "function_call" and item.get("call_id") == call_id: + item["status"] = "completed" + # Update arguments with parsed/sanitized version + item["arguments"] = tc.get("function", {}).get("arguments", "{}") + break + + for result in results: + output.append({ + "type": "function_call_output", + "id": output_id("fco"), + "call_id": result.get("tool_call_id", ""), + "output": [ + { + "type": "input_text", + "text": result.get("content", ""), + } + ], + "status": "completed", + **( + {"files": result.get("files")} + if result.get("files") + else {} + ), + **( + {"embeds": result.get("embeds")} + if result.get("embeds") + else {} + ), + }) + + # Append a new empty message item for the next response + output.append( { - "type": "text", - "content": "", + "type": "message", + "id": output_id("msg"), + "status": "in_progress", + "role": "assistant", + "content": [{"type": "output_text", "text": ""}], } ) @@ -4109,7 +3986,6 @@ async def streaming_chat_response_handler(response, ctx): ) tool_call_sources.clear() - output = convert_content_blocks_to_output(content_blocks) await event_emitter( { "type": "chat:completion", @@ -4127,9 +4003,7 @@ async def streaming_chat_response_handler(response, ctx): "stream": True, "messages": [ *form_data["messages"], - *convert_content_blocks_to_messages( - content_blocks, True - ), + *convert_output_to_messages(output, raw=True), ], } @@ -4153,11 +4027,11 @@ async def streaming_chat_response_handler(response, ctx): retries = 0 while ( - content_blocks[-1]["type"] == "code_interpreter" + output + and output[-1].get("type") == "open_webui:code_interpreter" and retries < MAX_RETRIES ): - output = convert_content_blocks_to_output(content_blocks) await event_emitter( { "type": "chat:completion", @@ -4171,10 +4045,11 @@ async def streaming_chat_response_handler(response, ctx): retries += 1 log.debug(f"Attempt count: {retries}") - output = "" + ci_item = output[-1] + ci_output = "" try: - if content_blocks[-1]["attributes"].get("type") == "code": - code = content_blocks[-1]["content"] + if ci_item.get("attributes", {}).get("type") == "code": + code = ci_item.get("code", "") # Sanitize code (strips ANSI codes and markdown fences) code = sanitize_code(code) @@ -4204,7 +4079,7 @@ async def streaming_chat_response_handler(response, ctx): request.app.state.config.CODE_INTERPRETER_ENGINE == "pyodide" ): - output = await event_caller( + ci_output = await event_caller( { "type": "execute:python", "data": { @@ -4220,7 +4095,7 @@ async def streaming_chat_response_handler(response, ctx): request.app.state.config.CODE_INTERPRETER_ENGINE == "jupyter" ): - output = await execute_code_jupyter( + ci_output = await execute_code_jupyter( request.app.state.config.CODE_INTERPRETER_JUPYTER_URL, code, ( @@ -4238,14 +4113,14 @@ async def streaming_chat_response_handler(response, ctx): request.app.state.config.CODE_INTERPRETER_JUPYTER_TIMEOUT, ) else: - output = { + ci_output = { "stdout": "Code interpreter engine not configured." } - log.debug(f"Code interpreter output: {output}") + log.debug(f"Code interpreter output: {ci_output}") - if isinstance(output, dict): - stdout = output.get("stdout", "") + if isinstance(ci_output, dict): + stdout = ci_output.get("stdout", "") if isinstance(stdout, str): stdoutLines = stdout.split("\n") @@ -4263,9 +4138,9 @@ async def streaming_chat_response_handler(response, ctx): f"![Output Image]({image_url})" ) - output["stdout"] = "\n".join(stdoutLines) + ci_output["stdout"] = "\n".join(stdoutLines) - result = output.get("result", "") + result = ci_output.get("result", "") if isinstance(result, str): resultLines = result.split("\n") @@ -4280,20 +4155,23 @@ async def streaming_chat_response_handler(response, ctx): resultLines[idx] = ( f"![Output Image]({image_url})" ) - output["result"] = "\n".join(resultLines) + ci_output["result"] = "\n".join(resultLines) except Exception as e: - output = str(e) + ci_output = str(e) - content_blocks[-1]["output"] = output + ci_item["output"] = ci_output + ci_item["status"] = "completed" - content_blocks.append( + output.append( { - "type": "text", - "content": "", + "type": "message", + "id": output_id("msg"), + "status": "in_progress", + "role": "assistant", + "content": [{"type": "output_text", "text": ""}], } ) - output = convert_content_blocks_to_output(content_blocks) await event_emitter( { "type": "chat:completion", @@ -4311,12 +4189,7 @@ async def streaming_chat_response_handler(response, ctx): "stream": True, "messages": [ *form_data["messages"], - { - "role": "assistant", - "content": serialize_content_blocks( - content_blocks, raw=True - ), - }, + *convert_output_to_messages(output, raw=True), ], } @@ -4335,8 +4208,12 @@ async def streaming_chat_response_handler(response, ctx): log.debug(e) break + # Mark all in-progress items as completed + for item in output: + if item.get("status") == "in_progress": + item["status"] = "completed" + title = Chats.get_chat_title_by_id(metadata["chat_id"]) - output = convert_content_blocks_to_output(content_blocks) data = { "done": True, "content": serialize_output(output), @@ -4392,7 +4269,6 @@ async def streaming_chat_response_handler(response, ctx): if not ENABLE_REALTIME_CHAT_SAVE: # Save message in the database - output = convert_content_blocks_to_output(content_blocks) Chats.upsert_message_to_chat_by_id_and_message_id( metadata["chat_id"], metadata["message_id"], diff --git a/backend/open_webui/utils/misc.py b/backend/open_webui/utils/misc.py index b931476ca..b2b10bf56 100644 --- a/backend/open_webui/utils/misc.py +++ b/backend/open_webui/utils/misc.py @@ -128,22 +128,40 @@ def get_content_from_message(message: dict) -> Optional[str]: return None -def convert_output_to_messages(output: list) -> list[dict]: +def convert_output_to_messages(output: list, raw: bool = False) -> list[dict]: """ - Convert OR-aligned output items to OpenAI-format messages for LLM consumption. - - This is the inverse of convert_content_blocks_to_output() in middleware.py. + Convert OR-aligned output items to OpenAI Chat Completion-format messages. + + This reconstructs the full conversation from the stored Responses API-native + output items, including assistant messages with tool_calls arrays and tool + role messages. + + Args: + output: List of OR-aligned output items (Responses API format). + raw: If True, include reasoning blocks (with original tags) and code + interpreter blocks for LLM re-processing follow-ups. """ if not output or not isinstance(output, list): return [] - + messages = [] pending_tool_calls = [] pending_content = [] - + + def flush_pending(): + nonlocal pending_content, pending_tool_calls + if pending_content or pending_tool_calls: + messages.append({ + "role": "assistant", + "content": "\n".join(pending_content) if pending_content else "", + **({"tool_calls": pending_tool_calls} if pending_tool_calls else {}), + }) + pending_content = [] + pending_tool_calls = [] + for item in output: item_type = item.get("type", "") - + if item_type == "message": # Extract text from output_text content parts content_parts = item.get("content", []) @@ -153,58 +171,86 @@ def convert_output_to_messages(output: list) -> list[dict]: text += part.get("text", "") if text: pending_content.append(text) - + elif item_type == "function_call": # Collect tool calls to batch into assistant message + arguments = item.get("arguments", "{}") + # Ensure arguments is always a JSON string + if not isinstance(arguments, str): + arguments = json.dumps(arguments) pending_tool_calls.append({ "id": item.get("call_id", ""), "type": "function", "function": { "name": item.get("name", ""), - "arguments": item.get("arguments", "{}"), + "arguments": arguments, } }) - + elif item_type == "function_call_output": # Flush any pending content/tool_calls before adding tool result - if pending_content or pending_tool_calls: - messages.append({ - "role": "assistant", - "content": "\n".join(pending_content) if pending_content else "", - **({"tool_calls": pending_tool_calls} if pending_tool_calls else {}), - }) - pending_content = [] - pending_tool_calls = [] - + flush_pending() + # Extract text from output content parts output_parts = item.get("output", []) content = "" for part in output_parts: if part.get("type") == "input_text": content += part.get("text", "") - + messages.append({ "role": "tool", "tool_call_id": item.get("call_id", ""), "content": content, }) - + elif item_type == "reasoning": - # Skip reasoning blocks for LLM messages - pass - + if raw: + # Include reasoning with original tags for LLM re-processing + reasoning_text = "" + source_list = item.get("summary", []) or item.get("content", []) + for part in source_list: + if part.get("type") == "output_text": + reasoning_text += part.get("text", "") + elif "text" in part: + reasoning_text += part.get("text", "") + + if reasoning_text: + start_tag = item.get("start_tag", "") + end_tag = item.get("end_tag", "") + pending_content.append( + f"{start_tag}{reasoning_text}{end_tag}" + ) + # 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", "") + + if code: + lang = item.get("lang", "python") + pending_content.append(f"```{lang}\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 + elif item_type.startswith("open_webui:"): - # Skip extension types + # Skip other extension types pass - + # Flush remaining content/tool_calls - if pending_content or pending_tool_calls: - messages.append({ - "role": "assistant", - "content": "\n".join(pending_content) if pending_content else "", - **({"tool_calls": pending_tool_calls} if pending_tool_calls else {}), - }) - + flush_pending() + return messages