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["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["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