Coverage for open_webui/utils/middleware.py: 3%
2896 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
1import ast
2import asyncio
3import base64
4import copy
5import html
6import inspect
7import json
8import logging
9import mimetypes
10import os
11import random
12import re
13import sys
14import textwrap
15import time
16from concurrent.futures import ThreadPoolExecutor
17from typing import Any, Optional
18from urllib.parse import unquote
19from uuid import uuid4
21from aiocache import cached
22from fastapi import HTTPException, Request
23from fastapi.responses import HTMLResponse, JSONResponse
24from open_webui.config import (
25 CACHE_DIR,
26 CODE_INTERPRETER_BLOCKED_MODULES,
27 CODE_INTERPRETER_PYODIDE_PROMPT,
28 DEFAULT_CODE_INTERPRETER_PROMPT,
29 DEFAULT_TOOLS_FUNCTION_CALLING_PROMPT_TEMPLATE,
30 DEFAULT_VOICE_MODE_PROMPT_TEMPLATE,
31)
32from open_webui.constants import TASKS
33from open_webui.env import (
34 BYPASS_MODEL_ACCESS_CONTROL,
35 CHAT_RESPONSE_MAX_TOOL_CALL_ITERATIONS,
36 CHAT_RESPONSE_STREAM_DELTA_CHUNK_SIZE,
37 ENABLE_API_OUTLET_FILTERS,
38 ENABLE_CHAT_RESPONSE_BASE64_IMAGE_URL_CONVERSION,
39 ENABLE_CHAT_RESPONSE_STREAM_INPLACE_APPEND,
40 ENABLE_PLUGINS,
41 ENABLE_QUERIES_CACHE,
42 ENABLE_REALTIME_CHAT_SAVE,
43 ENABLE_RESPONSES_API_STATEFUL,
44 GLOBAL_LOG_LEVEL,
45 RAG_SYSTEM_CONTEXT,
46)
47from open_webui.events import EVENTS, publish_event
48from open_webui.models.access_grants import AccessGrants
49from open_webui.models.chats import Chats
50from open_webui.models.config import Config
51from open_webui.models.folders import Folders
52from open_webui.models.models import Models
53from open_webui.models.notes import Notes
54from open_webui.models.oauth_sessions import OAuthSessions
55from open_webui.models.users import UserModel, Users
56from open_webui.retrieval.utils import filter_source_metadata, get_sources_from_items
57from open_webui.routers.images import (
58 CreateImageForm,
59 EditImageForm,
60 image_edits,
61 image_generations,
62)
63from open_webui.routers.pipelines import (
64 get_sorted_filters,
65 process_pipeline_inlet_filter,
66 process_pipeline_outlet_filter,
67)
68from open_webui.routers.retrieval import (
69 SearchForm,
70 process_web_search,
71)
72from open_webui.routers.tasks import (
73 generate_chat_tags,
74 generate_follow_ups,
75 generate_image_prompt,
76 generate_queries,
77 generate_title,
78)
79from open_webui.socket.main import (
80 get_event_call,
81 get_event_emitter,
82)
83from open_webui.tasks import clear_response_stream, save_response_stream
84from open_webui.utils.access_control import has_connection_access, has_permission
85from open_webui.utils.access_control.files import get_owner_accessible_folder_files
86from open_webui.utils.access_control.folders import has_folder_access
87from open_webui.utils.ask_user import stage_ask_user_tool_calls
88from open_webui.utils.chat import generate_chat_completion
89from open_webui.utils.chat_id import is_saved_chat_id
90from open_webui.utils.code_interpreter import execute_code_jupyter
91from open_webui.utils.context_compaction import compact_messages_for_request
92from open_webui.utils.files import (
93 convert_markdown_base64_images,
94 get_file_url_from_base64,
95 get_image_base64_from_url,
96 get_image_url_from_base64,
97)
98from open_webui.utils.filter import (
99 FilterContext,
100 get_filter_context,
101 get_filter_functions,
102 process_filter_functions,
103)
104from open_webui.utils.json_codec import JSONCodec
105from open_webui.utils.mcp.client import MCPClient
106from open_webui.utils.memory import add_memory_context, review_memory_after_turn
107from open_webui.utils.misc import (
108 add_or_update_system_message,
109 add_or_update_user_message,
110 convert_output_to_messages,
111 extract_urls,
112 get_content_from_message,
113 get_last_assistant_message,
114 get_last_user_message,
115 get_last_user_message_item,
116 get_message_list,
117 get_output_text,
118 get_response_error_detail,
119 get_reasoning_details,
120 get_system_message,
121 is_raster_image_content_type,
122 is_string_allowed,
123 merge_system_messages,
124 prepend_to_first_user_message_content,
125 replace_system_message_content,
126 set_last_user_message_content,
127 strip_empty_content_blocks,
128)
129from open_webui.utils.payload import apply_params_to_form_data, apply_system_prompt_to_body, resolve_system_prompt
130from open_webui.utils.plugin import load_function_module_by_id
131from open_webui.utils.response import merge_usage, normalize_usage
132from open_webui.utils.sanitize import sanitize_code
133from open_webui.utils.skills import (
134 apply_skills_create_prompt,
135 extract_skill_ids_from_messages,
136 has_prior_real_chat_content,
137 strip_skill_mentions,
138)
139from open_webui.utils.task import (
140 get_task_model_id,
141 rag_template,
142 tools_function_calling_generation_template,
143)
144from open_webui.utils.tools import (
145 build_tool_server_headers,
146 get_attached_knowledge,
147 get_builtin_tools,
148 get_terminal_tools,
149 get_tools,
150 get_updated_tool_function,
151)
152from starlette.responses import JSONResponse, Response, StreamingResponse
154logging.basicConfig(stream=sys.stdout, level=GLOBAL_LOG_LEVEL)
155log = logging.getLogger(__name__)
158def _is_tool_result_error(value: Any) -> bool:
159 if isinstance(value, str):
160 text = value.strip().lower()
161 if (
162 text.startswith('error:')
163 or text.startswith('exception:')
164 or text.startswith('traceback')
165 or text.startswith('http error!')
166 ):
167 return True
169 parsed = value
170 while isinstance(parsed, str):
171 try:
172 parsed = JSONCodec.loads(parsed)
173 except (JSONCodec.JSONDecodeError, TypeError, ValueError):
174 break
176 if not isinstance(parsed, dict):
177 return False
179 error = parsed.get('error')
180 if isinstance(error, str):
181 has_error = bool(error.strip())
182 else:
183 has_error = isinstance(error, (dict, list)) and bool(error)
184 if has_error:
185 return True
187 status = parsed.get('status')
188 if isinstance(status, str) and status.strip().lower() in {'error', 'failed'}:
189 return True
191 if parsed.get('success') is False or parsed.get('ok') is False:
192 message = parsed.get('message')
193 return has_error or (
194 bool(message.strip()) if isinstance(message, str) else isinstance(message, (dict, list)) and bool(message)
195 )
197 return False
200def normalize_messages_for_model(form_data: dict) -> dict:
201 form_data['messages'] = strip_empty_content_blocks(form_data.get('messages', []))
202 form_data['messages'] = merge_system_messages(form_data.get('messages', []))
203 return form_data
206async def publish_chat_finished_event(
207 request: Request, user: UserModel, metadata: dict, title: str, content: str, output: list | None = None
208):
209 chat_id = metadata.get('chat_id')
210 if getattr(request.state, 'internal', False) is True or not is_saved_chat_id(chat_id):
211 return
213 content = content or get_output_text(output)
214 webui_url = await Config.get('webui.url')
215 await publish_event(
216 request,
217 EVENTS.CHAT_FINISHED,
218 actor=user,
219 subject_id=chat_id,
220 subject_type='chat',
221 data={
222 'user_id': user.id,
223 'chat_id': chat_id,
224 'message_id': metadata.get('message_id'),
225 'model_id': metadata.get('model_id'),
226 'title': title,
227 'url': f'{webui_url}/c/{chat_id}' if webui_url else f'/c/{chat_id}',
228 'message': content,
229 },
230 message=title or 'Chat finished',
231 )
232 event_emitter = await get_event_emitter(metadata, update_db=False)
233 if event_emitter:
234 folder_id = metadata.get('folder_id') or await Chats.get_chat_folder_id(chat_id, metadata.get('user_id'))
235 await event_emitter({'type': 'chat:list', 'data': {'chat_id': chat_id, 'folder_id': folder_id}})
238# We believe in one maker of all models, seen and unseen,
239# and in the reasoning which proceeds from the architect.
240# We look for the resurrection of dead processes and the
241# inference of the world to come.
242DEFAULT_REASONING_TAGS = [
243 ('<think>', '</think>'),
244 ('<thinking>', '</thinking>'),
245 ('<reason>', '</reason>'),
246 ('<reasoning>', '</reasoning>'),
247 ('<thought>', '</thought>'),
248 ('<Thought>', '</Thought>'),
249 ('<|begin_of_thought|>', '<|end_of_thought|>'),
250 ('◁think▷', '◁/think▷'),
251]
253DEFAULT_SOLUTION_TAGS = [('<|begin_of_solution|>', '<|end_of_solution|>')]
254DEFAULT_CODE_INTERPRETER_TAGS = [('<code_interpreter>', '</code_interpreter>')]
257def _start_tag_pattern(start_tag: str) -> str:
258 if start_tag.startswith('<') and start_tag.endswith('>'):
259 return rf'<{re.escape(start_tag[1:-1])}(\s.*?)?>'
260 return re.escape(start_tag)
263def output_id(prefix: str) -> str:
264 """Generate OR-style ID: prefix + 24-char hex UUID."""
265 return f'{prefix}_{uuid4().hex[:24]}'
268def build_terminal_file_tool_result(
269 tool_function_name: str,
270 tool_function_params: dict,
271 tool_result: Any,
272 tool: dict | None,
273 metadata: dict | None,
274) -> dict | None:
275 if isinstance(tool_result, (list, tuple)) and tool_result and isinstance(tool_result[0], dict):
276 tool_result = tool_result[0]
278 if tool_function_name != 'display_file' or not isinstance(tool_result, dict) or tool_result.get('exists') is False:
279 return None
281 tool_id = (tool or {}).get('tool_id', '')
282 terminal_id = metadata.get('terminal_id') if metadata else None
283 if isinstance(tool_id, str) and tool_id.startswith('terminal:'):
284 terminal_id = tool_id.split(':', 1)[1]
286 server_url = ((tool or {}).get('server') or {}).get('url')
287 terminal_selector = terminal_id or server_url
288 path = tool_result.get('path') or tool_function_params.get('path')
289 if not terminal_selector or not path:
290 return None
291 mime_type, _ = mimetypes.guess_type(path)
292 mime_type = mime_type or 'application/octet-stream'
293 page = tool_result.get('page') or tool_function_params.get('page')
295 return {
296 **tool_result,
297 'type': 'file',
298 'source': 'open_terminal',
299 **({'displayed': True} if tool_function_params.get('inline') is True else {}),
300 'terminal_selector': terminal_selector,
301 **({'terminal_id': terminal_id} if terminal_id else {}),
302 **({'terminal_url': server_url} if server_url and not terminal_id else {}),
303 'session_id': metadata.get('chat_id') if metadata else None,
304 'path': path,
305 'full_path': tool_result.get('full_path') or path,
306 'name': tool_result.get('name') or os.path.basename(path),
307 'mime_type': tool_result.get('mime_type') or tool_result.get('content_type') or mime_type,
308 'content_type': tool_result.get('content_type') or tool_result.get('mime_type') or mime_type,
309 **({'page': page} if page else {}),
310 }
313def tool_result_content(tool_result: Any) -> str:
314 if not tool_result:
315 return ''
316 if isinstance(tool_result, (dict, list)):
317 return JSONCodec.dumps(tool_result, ensure_ascii=False)
318 return str(tool_result)
321def append_to_text_field(item: dict, key: str, value: str) -> None:
322 # Default: keep the existing field intact until concatenation succeeds.
323 # Non-string values and str subclasses also use this path to preserve errors.
324 if not ENABLE_CHAT_RESPONSE_STREAM_INPLACE_APPEND or type(item[key]) is not str or type(value) is not str:
325 item[key] += value
326 return
328 # Opt-in: dropping the dict's reference lets CPython extend an unshared str
329 # in place. Only a host that is already out of memory can leave the field empty.
330 text = item[key]
331 item[key] = ''
332 text += value
333 item[key] = text
336def merge_streamed_reasoning_details(target: list, details) -> None:
337 items = details if isinstance(details, list) else [details]
338 for item in items:
339 if not isinstance(item, dict):
340 continue
342 index = item.get('index')
343 existing = (
344 next((detail for detail in target if detail.get('index') == index), None)
345 if isinstance(index, int)
346 else None
347 )
348 if existing is None:
349 target.append(dict(item))
350 continue
352 for key, value in item.items():
353 if key in ('text', 'summary') and isinstance(value, str) and isinstance(existing.get(key), str):
354 append_to_text_field(existing, key, value)
355 else:
356 existing[key] = value
359def _split_tool_calls(
360 tool_calls: list[dict],
361) -> list[dict]:
362 """Expand tool calls whose arguments contain multiple back-to-back JSON objects.
364 Some models (e.g. GPT-5.4) send multiple complete JSON argument objects
365 under the same tool call index, producing concatenated invalid JSON like:
366 '{"query":"A","count":5}{"query":"B","count":5}'
368 Each such tool call is split into separate entries so each gets executed
369 independently. Single-object arguments pass through unchanged.
370 """
372 def split_json_objects(raw: str) -> list[str]:
373 if not isinstance(raw, str):
374 raw = '' if raw is None else JSONCodec.dumps(raw)
376 decoder = json.JSONDecoder()
377 results = []
378 position = 0
380 while position < len(raw):
381 while position < len(raw) and raw[position].isspace():
382 position += 1
383 if position >= len(raw):
384 break
385 try:
386 _, end = decoder.raw_decode(raw, position)
387 results.append(raw[position:end].strip())
388 position = end
389 except JSONCodec.JSONDecodeError:
390 return [raw]
392 return results or [raw]
394 expanded = []
395 for tool_call in tool_calls:
396 function = tool_call.setdefault('function', {})
397 arguments = function.get('arguments')
398 if not isinstance(arguments, str):
399 arguments = '' if arguments is None else JSONCodec.dumps(arguments)
400 function['arguments'] = arguments
401 split_arguments = split_json_objects(arguments)
403 if len(split_arguments) <= 1:
404 expanded.append(tool_call)
405 else:
406 for argument in split_arguments:
407 cloned = copy.deepcopy(tool_call)
408 cloned['id'] = f'call_{uuid4().hex[:24]}'
409 cloned['function']['arguments'] = argument
410 expanded.append(cloned)
412 return expanded
415def get_citation_source_from_tool_result(
416 tool_name: str, tool_params: dict, tool_result: str, tool_id: str = ''
417) -> list[dict]:
418 """
419 Parse a tool's result and convert it to source dicts for citation display.
421 Follows the source format conventions from get_sources_from_items:
422 - source: file/item info object with id, name, type
423 - document: list of document contents
424 - metadata: list of metadata objects with source, file_id, name fields
426 Returns a list of sources (usually one, but query_knowledge_files/query_chat_files may return multiple).
427 """
428 try:
429 try:
430 tool_result = JSONCodec.loads(tool_result)
431 except (JSONCodec.JSONDecodeError, TypeError):
432 pass # keep tool_result as-is (e.g. fetch_url returns plain text)
433 if isinstance(tool_result, dict) and 'error' in tool_result:
434 return []
436 if tool_name in ('view_knowledge_file', 'view_file'):
437 if not isinstance(tool_result, dict):
438 return []
440 file_data = tool_result
441 filename = file_data.get('filename', 'Unknown File')
442 file_id = file_data.get('id', '')
443 knowledge_name = file_data.get('knowledge_name', '')
445 return [
446 {
447 'source': {
448 'id': file_id,
449 'name': filename,
450 'type': 'file',
451 },
452 'document': [file_data.get('content', '')],
453 'metadata': [
454 {
455 'file_id': file_id,
456 'name': filename,
457 'source': filename,
458 **({'knowledge_name': knowledge_name} if knowledge_name else {}),
459 }
460 ],
461 }
462 ]
464 elif tool_name == 'fetch_url':
465 url = tool_params.get('url', '')
466 content = tool_result if isinstance(tool_result, str) else str(tool_result)
467 snippet = content[:500] + ('...' if len(content) > 500 else '')
469 return [
470 {
471 'source': {'name': url or 'fetch_url', 'id': url or 'fetch_url'},
472 'document': [snippet],
473 'metadata': [
474 {
475 'source': url,
476 'name': url,
477 'url': url,
478 }
479 ],
480 }
481 ]
483 elif tool_name in ('query_knowledge_files', 'query_chat_files'):
484 if not isinstance(tool_result, list):
485 return []
487 chunks = tool_result
489 # Group chunks by source for better citation display
490 # Each unique source becomes a separate source entry
491 sources_by_file = {}
493 for chunk in chunks:
494 source_name = chunk.get('source', 'Unknown')
495 file_id = chunk.get('file_id', '')
496 note_id = chunk.get('note_id', '')
497 chunk_type = chunk.get('type', 'file')
498 content = chunk.get('content', '')
500 # Use file_id or note_id as the key
501 key = file_id or note_id or source_name
503 if key not in sources_by_file:
504 sources_by_file[key] = {
505 'source': {
506 'id': file_id or note_id,
507 'name': source_name,
508 'type': chunk_type,
509 },
510 'document': [],
511 'metadata': [],
512 }
514 sources_by_file[key]['document'].append(content)
515 sources_by_file[key]['metadata'].append(
516 {
517 'file_id': file_id,
518 'name': source_name,
519 'source': source_name,
520 **({'note_id': note_id} if note_id else {}),
521 }
522 )
524 # Return all grouped sources as a list
525 if sources_by_file:
526 return list(sources_by_file.values())
528 # Empty result fallback
529 return []
531 else:
532 # Fallback for other tools
533 return [
534 {
535 'source': {
536 'name': tool_name,
537 'type': 'tool',
538 'id': tool_id or tool_name,
539 },
540 'document': [str(tool_result)],
541 'metadata': [{'source': tool_name, 'name': tool_name}],
542 }
543 ]
544 except Exception as e:
545 log.exception(f'Error parsing tool result for {tool_name}: {e}')
546 return [
547 {
548 'source': {'name': tool_name, 'type': 'tool'},
549 'document': [str(tool_result)],
550 'metadata': [{'source': tool_name}],
551 }
552 ]
555def deep_merge(target, source):
556 """
557 Merge source into target recursively (returning new structure).
558 - Dicts: Recursive merge.
559 - Strings: Concatenation.
560 - Others: Overwrite.
561 """
562 if isinstance(target, dict) and isinstance(source, dict):
563 new_target = target.copy()
564 for k, v in source.items():
565 if k in new_target:
566 new_target[k] = deep_merge(new_target[k], v)
567 else:
568 new_target[k] = v
569 return new_target
570 elif isinstance(target, str) and isinstance(source, str):
571 return target + source
572 else:
573 return source
576RESPONSE_COMPLETION_RESPONSE_FIELDS = ('error', 'id', 'output', 'usage')
579def get_response_completion_event_data(event: dict) -> dict:
580 """Build the data payload for response:completion events."""
581 response = event.get('response')
582 if not isinstance(response, dict):
583 return event
585 response_data = {key: response[key] for key in RESPONSE_COMPLETION_RESPONSE_FIELDS if key in response}
587 return {
588 **event,
589 'response': response_data,
590 }
593def handle_responses_streaming_event(
594 data: dict,
595 current_output: list,
596) -> tuple[list, dict | None]:
597 """
598 Handle Responses API streaming events in a pure functional way.
600 Args:
601 data: The event data
602 current_output: List of output items (treated as immutable)
604 Returns:
605 tuple[list, dict | None]: (new_output, metadata)
606 - new_output: The updated output list.
607 - metadata: Metadata to emit (e.g. usage), {} if update occurred, None if skip.
608 """
609 # Default: no change
610 # Note: treating current_output as immutable, but avoiding full deepcopy for perf.
611 # We will shallow copy only if we need to modify the list structure or items.
613 event_type = data.get('type', '')
615 if event_type == 'response.output_item.added':
616 item = data.get('item', {})
617 if item:
618 new_output = list(current_output)
619 output_index = data.get('output_index', len(new_output))
620 existing_index = next(
621 (
622 idx
623 for idx, existing in enumerate(new_output)
624 if (item.get('id') and existing.get('id') == item.get('id'))
625 or (
626 item.get('call_id')
627 and existing.get('type') == item.get('type')
628 and existing.get('call_id') == item.get('call_id')
629 )
630 ),
631 None,
632 )
633 if existing_index is not None:
634 new_output[existing_index] = item
635 elif 0 <= output_index < len(new_output):
636 new_output.insert(output_index, item)
637 else:
638 new_output.append(item)
639 return new_output, None
640 return current_output, None
642 elif event_type == 'response.content_part.added':
643 part = data.get('part', {})
644 output_index = data.get('output_index', len(current_output) - 1)
646 if current_output and 0 <= output_index < len(current_output):
647 new_output = list(current_output)
648 # Copy the item to mutate it
649 item = new_output[output_index].copy()
650 new_output[output_index] = item
652 if 'content' not in item:
653 item['content'] = []
654 else:
655 # Copy content list
656 item['content'] = list(item['content'])
658 if item.get('type') == 'reasoning':
659 # Reasoning items should not have content parts
660 pass
661 else:
662 item['content'].append(part)
663 return new_output, None
664 return current_output, None
666 elif event_type == 'response.reasoning_summary_part.added':
667 part = data.get('part', {})
668 output_index = data.get('output_index', len(current_output) - 1)
670 if current_output and 0 <= output_index < len(current_output):
671 new_output = list(current_output)
672 item = new_output[output_index].copy()
673 new_output[output_index] = item
675 if 'summary' not in item:
676 item['summary'] = []
677 else:
678 item['summary'] = list(item['summary'])
680 item['summary'].append(part)
681 return new_output, None
682 return current_output, None
684 elif event_type.startswith('response.') and event_type.endswith('.delta'):
685 # Generic Delta Handling
686 parts = event_type.split('.')
687 if len(parts) >= 3:
688 delta_type = parts[1]
689 delta = data.get('delta', '')
691 output_index = data.get('output_index', len(current_output) - 1)
693 if current_output and 0 <= output_index < len(current_output):
694 new_output = list(current_output)
695 item = new_output[output_index].copy()
696 new_output[output_index] = item
697 item_type = item.get('type', '')
699 # Determine target field and object based on delta_type and item_type
700 if delta_type == 'function_call_arguments':
701 key = 'arguments'
702 if item_type == 'function_call':
703 # Function call args are usually strings
704 item[key] = item.get(key, '') + str(delta)
705 else:
706 # Generic handling, refined by item type below
707 pass
709 if item_type == 'message':
710 # Message items: "text"/"output_text" -> "text"
711 # "reasoning_text" -> Skipped (should use reasoning item)
712 if delta_type in ['text', 'output_text']:
713 key = 'text'
714 elif delta_type in ['reasoning_text', 'reasoning_summary_text']:
715 # Skip reasoning updates for message items
716 return new_output, None
717 else:
718 key = delta_type
720 content_index = data.get('content_index', 0)
721 if 'content' not in item:
722 item['content'] = []
723 else:
724 item['content'] = list(item['content'])
725 content_list = item['content']
727 while len(content_list) <= content_index:
728 content_list.append({'type': 'text', 'text': ''})
730 # Copy the part to mutate it
731 part = content_list[content_index].copy()
732 content_list[content_index] = part
734 current_val = part.get(key)
735 if current_val is None:
736 # Initialize based on delta type
737 current_val = {} if isinstance(delta, dict) else ''
739 part[key] = deep_merge(current_val, delta)
741 elif item_type == 'reasoning':
742 # Reasoning items: "reasoning_text"/"reasoning_summary_text" -> "text"
743 # "text"/"output_text" -> Skipped (should use message item)
744 if delta_type == 'reasoning_summary_text':
745 # Summary updates -> item['summary']
746 key = 'text'
747 summary_index = data.get('summary_index', 0)
748 if 'summary' not in item:
749 item['summary'] = []
750 else:
751 item['summary'] = list(item['summary'])
752 summary_list = item['summary']
754 while len(summary_list) <= summary_index:
755 summary_list.append({'type': 'summary_text', 'text': ''})
757 part = summary_list[summary_index].copy()
758 summary_list[summary_index] = part
760 target_val = part.get(key, '')
761 part[key] = deep_merge(target_val, delta)
763 elif delta_type == 'reasoning_text':
764 # Reasoning body updates -> item['content']
765 key = 'text'
766 content_index = data.get('content_index', 0)
767 if 'content' not in item:
768 item['content'] = []
769 else:
770 item['content'] = list(item['content'])
771 content_list = item['content']
773 while len(content_list) <= content_index:
774 # Reasoning content parts default to text
775 content_list.append({'type': 'text', 'text': ''})
777 part = content_list[content_index].copy()
778 content_list[content_index] = part
780 target_val = part.get(key, '')
781 part[key] = deep_merge(target_val, delta)
783 elif delta_type in ['text', 'output_text']:
784 return new_output, None
785 else:
786 # Fallback just in case other deltas target reasoning?
787 pass
789 else:
790 # Fallback for other item types
791 if delta_type in ['text', 'output_text']:
792 key = 'text'
793 else:
794 key = delta_type
796 current_val = item.get(key)
797 if current_val is None:
798 current_val = {} if isinstance(delta, dict) else ''
799 item[key] = deep_merge(current_val, delta)
801 return new_output, None
803 return current_output, None
805 elif event_type == 'response.output_item.done':
806 # Delta Event: Output item complete
807 item = data.get('item')
808 output_index = data.get('output_index', len(current_output) - 1)
810 new_output = list(current_output)
811 if item and 0 <= output_index < len(current_output):
812 new_output[output_index] = item
813 elif item:
814 new_output.append(item)
815 return new_output, {}
817 elif event_type.startswith('response.') and event_type.endswith('.done'):
818 # Delta Events: response.content_part.done, response.text.done, etc.
819 parts = event_type.split('.')
820 if len(parts) >= 3:
821 type_name = parts[1]
823 # 1. Handle specific Delta "done" signals
824 if type_name == 'content_part':
825 # "Signaling that no further changes will occur to a content part"
826 # If payloads contains the full part, we could update it.
827 # Usually purely signaling in standard implementation, but we check payload.
828 part = data.get('part')
829 output_index = data.get('output_index', len(current_output) - 1)
831 if part and current_output and 0 <= output_index < len(current_output):
832 new_output = list(current_output)
833 item = new_output[output_index].copy()
834 new_output[output_index] = item
836 if 'content' in item:
837 item['content'] = list(item['content'])
838 content_index = data.get('content_index', len(item['content']) - 1)
839 if 0 <= content_index < len(item['content']):
840 item['content'][content_index] = part
841 return new_output, {}
842 return current_output, None
844 elif type_name == 'reasoning_summary_part':
845 part = data.get('part')
846 output_index = data.get('output_index', len(current_output) - 1)
848 if part and current_output and 0 <= output_index < len(current_output):
849 new_output = list(current_output)
850 item = new_output[output_index].copy()
851 new_output[output_index] = item
853 if 'summary' in item:
854 item['summary'] = list(item['summary'])
855 summary_index = data.get('summary_index', len(item['summary']) - 1)
856 if 0 <= summary_index < len(item['summary']):
857 item['summary'][summary_index] = part
858 return new_output, {}
859 return current_output, None
861 # 2. Generic Field Done (text.done, audio.done)
862 if type_name not in ['completed', 'failed']:
863 output_index = data.get('output_index', len(current_output) - 1)
864 if current_output and 0 <= output_index < len(current_output):
865 key = (
866 'text'
867 if type_name
868 in [
869 'text',
870 'output_text',
871 'reasoning_text',
872 'reasoning_summary_text',
873 ]
874 else type_name
875 )
876 if type_name == 'function_call_arguments':
877 key = 'arguments'
879 if key in data:
880 final_value = data[key]
881 new_output = list(current_output)
882 item = new_output[output_index].copy()
883 new_output[output_index] = item
884 item_type = item.get('type', '')
886 if type_name == 'function_call_arguments':
887 if item_type == 'function_call':
888 item['arguments'] = final_value
889 elif item_type == 'message':
890 content_index = data.get('content_index', 0)
891 if 'content' in item:
892 item['content'] = list(item['content'])
893 if len(item['content']) > content_index:
894 part = item['content'][content_index].copy()
895 item['content'][content_index] = part
896 part[key] = final_value
897 elif item_type == 'reasoning':
898 item['status'] = 'completed'
899 else:
900 item[key] = final_value
902 return new_output, {}
904 return current_output, None
906 elif event_type == 'response.completed':
907 # State Machine Event: Completed
908 response_data = data.get('response', {})
909 final_output = response_data.get('output')
911 # Some providers send an empty output on response.completed despite having streamed items
912 new_output = final_output if final_output else current_output
914 # Ensure reasoning items are marked as completed in the final output
915 if new_output:
916 for item in new_output:
917 if item.get('type') == 'reasoning' and item.get('status') != 'completed':
918 item['status'] = 'completed'
920 return new_output, {
921 'usage': response_data.get('usage'),
922 'done': True,
923 'response_id': response_data.get('id'),
924 }
926 elif event_type == 'response.in_progress':
927 # State Machine Event: In Progress
928 # We could extract metadata if needed, but for now just acknowledge iteration
929 return current_output, None
931 elif event_type == 'response.failed':
932 # State Machine Event: Failed
933 error = data.get('response', {}).get('error', {})
934 return current_output, {'error': error}
936 else:
937 return current_output, None
940def get_source_context(sources: list, source_ids: dict = None, include_content: bool = True) -> str:
941 """
942 Build <source> tag context string from citation sources.
943 """
944 context_string = ''
945 if source_ids is None:
946 source_ids = {}
947 for source in sources:
948 for doc, meta in zip(source.get('document', []), source.get('metadata', [])):
949 source_id = meta.get('source') or source.get('source', {}).get('id') or 'N/A'
950 if source_id not in source_ids:
951 source_ids[source_id] = len(source_ids) + 1
952 src_name = source.get('source', {}).get('name')
953 src_type = source.get('source', {}).get('type')
954 src_rid = source.get('source', {}).get('id')
955 body = doc if include_content else ''
956 extra_attrs = ''
957 for key, value in filter_source_metadata(meta).items():
958 if key in ('id', 'name', 'resource-type', 'resource-id'):
959 continue
960 extra_attrs += f' {key}="{html.escape(str(value))}"'
961 context_string += (
962 f'<source id="{source_ids[source_id]}"'
963 + (f' name="{src_name}"' if src_name else '')
964 + (f' resource-type="{src_type}"' if src_type else '')
965 + (f' resource-id="{src_rid}"' if src_rid else '')
966 + extra_attrs
967 + f'>{body}</source>\n'
968 )
969 return context_string
972async def apply_source_context_to_messages(
973 request: Request,
974 messages: list,
975 sources: list,
976 user_message: str,
977 include_content: bool = True,
978) -> list:
979 """
980 Build source context from citation sources and apply to messages.
981 Uses RAG template to format context for model consumption.
983 When include_content is False, emit <source> tags with id/name but no
984 document body — useful when the content is already present elsewhere
985 (e.g. in a tool result message) and only citation markers are needed.
986 """
987 if not sources or not user_message:
988 return messages
990 context = get_source_context(sources, include_content=include_content)
992 context = context.strip()
993 if not context:
994 return messages
996 if RAG_SYSTEM_CONTEXT:
997 return add_or_update_system_message(
998 await rag_template(await Config.get('rag.template'), context, user_message),
999 messages,
1000 append=True,
1001 )
1002 else:
1003 return add_or_update_user_message(
1004 await rag_template(await Config.get('rag.template'), context, user_message),
1005 messages,
1006 append=False,
1007 )
1010BASE64_IMAGE_DATA_URI_RE = re.compile(r'data:image/[a-zA-Z0-9.+-]+;base64,[A-Za-z0-9+/]+={0,2}', re.IGNORECASE)
1013def extract_base64_images(value: Any, files: list) -> Any:
1014 """Move base64 image data URIs out of a tool result so they do not reach the model as text."""
1015 if isinstance(value, str):
1016 if BASE64_IMAGE_DATA_URI_RE.fullmatch(value):
1017 files.append({'type': 'image', 'url': value})
1018 return '[image]'
1019 return value
1020 if isinstance(value, dict):
1021 return {key: extract_base64_images(item, files) for key, item in value.items()}
1022 if isinstance(value, list):
1023 return [extract_base64_images(item, files) for item in value]
1024 if isinstance(value, tuple):
1025 return tuple(extract_base64_images(item, files) for item in value)
1026 return value
1029async def store_tool_result_image(request, image_url, metadata, user):
1030 """Keep saved tool images out of chat JSON, falling back to inline data if storage fails."""
1031 metadata = metadata or {}
1032 if (
1033 not isinstance(image_url, str)
1034 or not image_url.startswith('data:image/')
1035 or not is_saved_chat_id(metadata.get('chat_id'))
1036 ):
1037 return image_url
1039 try:
1040 stored_url = await get_file_url_from_base64(
1041 request,
1042 image_url,
1043 {key: metadata.get(key) for key in ('chat_id', 'message_id', 'session_id')},
1044 user,
1045 )
1046 return stored_url or image_url
1047 except Exception:
1048 log.warning('Could not store tool image; retaining inline image')
1049 return image_url
1052async def process_tool_result(
1053 request,
1054 tool_function_name,
1055 tool_result,
1056 tool_type,
1057 direct_tool=False,
1058 metadata=None,
1059 user=None,
1060):
1061 tool_result_embeds = []
1062 EXTERNAL_TOOL_TYPES = ('external', 'action', 'terminal')
1064 # Support (HTMLResponse, result_context) tuples: the optional second
1065 # element lets tool authors provide the LLM with actionable context
1066 # about the generated embed instead of the generic fallback message.
1067 result_context = None
1068 if isinstance(tool_result, tuple) and len(tool_result) == 2 and isinstance(tool_result[0], HTMLResponse):
1069 tool_result, result_context = tool_result
1071 if isinstance(tool_result, HTMLResponse):
1072 content_disposition = tool_result.headers.get('Content-Disposition', '')
1073 if 'inline' in content_disposition:
1074 content = tool_result.body.decode('utf-8', 'replace')
1075 tool_result_embeds.append(content)
1077 if 200 <= tool_result.status_code < 300:
1078 if result_context is not None and isinstance(result_context, (str, dict, list)):
1079 tool_result = result_context
1080 else:
1081 tool_result = {
1082 'status': 'success',
1083 'code': 'ui_component',
1084 'message': f'{tool_function_name}: Embedded UI result is active and visible to the user.',
1085 }
1086 elif 400 <= tool_result.status_code < 500:
1087 tool_result = {
1088 'status': 'error',
1089 'code': 'ui_component',
1090 'message': f'{tool_function_name}: Client error {tool_result.status_code} from embedded UI result.',
1091 }
1092 elif 500 <= tool_result.status_code < 600:
1093 tool_result = {
1094 'status': 'error',
1095 'code': 'ui_component',
1096 'message': f'{tool_function_name}: Server error {tool_result.status_code} from embedded UI result.',
1097 }
1098 else:
1099 tool_result = {
1100 'status': 'error',
1101 'code': 'ui_component',
1102 'message': f'{tool_function_name}: Unexpected status code {tool_result.status_code} from embedded UI result.',
1103 }
1104 else:
1105 tool_result = tool_result.body.decode('utf-8', 'replace')
1107 elif (tool_type in EXTERNAL_TOOL_TYPES and isinstance(tool_result, tuple)) or (
1108 direct_tool and isinstance(tool_result, list) and len(tool_result) == 2
1109 ):
1110 tool_result, tool_response_headers = tool_result
1112 try:
1113 if not isinstance(tool_response_headers, dict):
1114 tool_response_headers = dict(tool_response_headers)
1115 except Exception as e:
1116 tool_response_headers = {}
1117 log.debug(e)
1119 if tool_response_headers and isinstance(tool_response_headers, dict):
1120 content_disposition = tool_response_headers.get(
1121 'Content-Disposition',
1122 tool_response_headers.get('content-disposition', ''),
1123 )
1125 if 'inline' in content_disposition:
1126 content_type = tool_response_headers.get(
1127 'Content-Type',
1128 tool_response_headers.get('content-type', ''),
1129 )
1130 location = tool_response_headers.get(
1131 'Location',
1132 tool_response_headers.get('location', ''),
1133 )
1135 if 'text/html' in content_type:
1136 # Support (html_content, result_context) nested tuple
1137 result_context = None
1138 html_content = tool_result
1139 if isinstance(tool_result, (tuple, list)) and len(tool_result) == 2:
1140 html_content, result_context = tool_result
1142 # Display as iframe embed
1143 tool_result_embeds.append(html_content)
1144 if result_context is not None and isinstance(result_context, (str, dict, list)):
1145 tool_result = result_context
1146 else:
1147 tool_result = {
1148 'status': 'success',
1149 'code': 'ui_component',
1150 'message': f'{tool_function_name}: Embedded UI result is active and visible to the user.',
1151 }
1152 elif location:
1153 # Support (html_content, result_context) nested tuple for location embeds
1154 result_context = None
1155 if isinstance(tool_result, (tuple, list)) and len(tool_result) == 2:
1156 _, result_context = tool_result
1158 tool_result_embeds.append(location)
1159 if result_context is not None and isinstance(result_context, (str, dict, list)):
1160 tool_result = result_context
1161 else:
1162 tool_result = {
1163 'status': 'success',
1164 'code': 'ui_component',
1165 'message': f'{tool_function_name}: Embedded UI result is active and visible to the user.',
1166 }
1168 tool_result_files = []
1170 # Detect base64 image data URIs from tool results (e.g. binary image
1171 # responses from execute_tool_server). Move the data URI to
1172 # tool_result_files and replace tool_result with a text summary.
1173 if isinstance(tool_result, str) and tool_result.startswith('data:image/'):
1174 tool_result_files.append({'type': 'image', 'url': tool_result})
1175 tool_result = f'{tool_function_name}: Image file read successfully.'
1177 if isinstance(tool_result, list):
1178 if tool_type == 'mcp': # MCP
1179 tool_response = []
1180 for item in tool_result:
1181 if isinstance(item, dict):
1182 if item.get('type') == 'text':
1183 text = item.get('text', '')
1184 if isinstance(text, str):
1185 try:
1186 text = JSONCodec.loads(text)
1187 except JSONCodec.JSONDecodeError:
1188 pass
1189 tool_response.append(text)
1190 elif item.get('type') in ['image', 'audio']:
1191 file_url = await get_file_url_from_base64(
1192 request,
1193 f'data:{item.get("mimeType")};base64,{item.get("data", item.get("blob", ""))}',
1194 {
1195 'chat_id': metadata.get('chat_id', None),
1196 'message_id': metadata.get('message_id', None),
1197 'session_id': metadata.get('session_id', None),
1198 'result': item,
1199 },
1200 user,
1201 )
1203 tool_result_files.append(
1204 {
1205 'type': item.get('type', 'data'),
1206 'url': file_url,
1207 }
1208 )
1209 elif item.get('type') == 'resource':
1210 resource = item.get('resource', {})
1211 text = resource.get('text', '')
1212 if isinstance(text, str) and text:
1213 try:
1214 text = JSONCodec.loads(text)
1215 except JSONCodec.JSONDecodeError:
1216 pass
1217 tool_response.append(text)
1218 elif resource.get('blob'):
1219 resource_mime_type = resource.get('mimeType') or 'application/octet-stream'
1220 resource_blob = resource.get('blob', '')
1221 if resource_mime_type.startswith('image/'):
1222 tool_result_files.append(
1223 {
1224 'type': 'image',
1225 'url': f'data:{resource_mime_type};base64,{resource_blob}',
1226 }
1227 )
1228 else:
1229 resource_uri = resource.get('uri', 'resource')
1230 tool_response.append(
1231 f'[Resource: {resource_uri}] (binary data, mimeType: {resource_mime_type})'
1232 )
1233 elif resource.get('uri'):
1234 tool_response.append(resource.get('uri'))
1235 tool_result = tool_response[0] if len(tool_response) == 1 else tool_response
1236 else: # OpenAPI
1237 # Images are left to extract_base64_images below, which attaches them so the model can see them.
1238 for item in list(tool_result):
1239 if isinstance(item, str) and item.startswith('data:') and not BASE64_IMAGE_DATA_URI_RE.fullmatch(item):
1240 tool_result_files.append(
1241 {
1242 'type': 'data',
1243 'content': item,
1244 }
1245 )
1246 tool_result.remove(item)
1248 tool_result = extract_base64_images(tool_result, tool_result_files)
1250 if isinstance(tool_result, list):
1251 tool_result = {'results': tool_result}
1253 if isinstance(tool_result, dict) or isinstance(tool_result, list):
1254 tool_result = json.dumps(tool_result, indent=2, ensure_ascii=False)
1256 # Safety: ensure tool_result is always a string (or None) to prevent
1257 # downstream TypeError when concatenating (e.g. if an upstream callable
1258 # returned a tuple that was not unpacked by the branches above).
1259 if tool_result is not None and not isinstance(tool_result, str):
1260 if isinstance(tool_result, tuple):
1261 # execute_tool_server returns (data, headers); unpack the data part
1262 tool_result = json.dumps(tool_result[0], indent=2, ensure_ascii=False) if len(tool_result) > 0 else ''
1263 else:
1264 tool_result = str(tool_result)
1266 return tool_result, tool_result_files, tool_result_embeds
1269def parse_terminal_tool_result(tool_result: Any) -> dict:
1270 try:
1271 result = JSONCodec.loads(tool_result)
1272 except (JSONCodec.JSONDecodeError, TypeError):
1273 return {}
1274 return result if isinstance(result, dict) else {}
1277async def terminal_event_handler(
1278 tool_function_name: str,
1279 tool_function_params: dict,
1280 tool_result,
1281 event_emitter,
1282):
1283 """Emit terminal:* events for Open Terminal tools.
1285 - display_file → emits 'terminal:display_file' to open the file preview.
1286 - write_file / replace_file_content → emits 'terminal:write_file' to refresh.
1287 - run_command → emits 'terminal:run_command' with cwd to refresh if relevant.
1288 """
1289 if not event_emitter:
1290 return
1292 if tool_function_name == 'display_file':
1293 if tool_function_params.get('inline') is True:
1294 return
1295 result = parse_terminal_tool_result(tool_result)
1296 # Open Terminal resolves the argument against the session cwd, so prefer its resolved path
1297 path = result.get('path') or tool_function_params.get('path', '')
1298 if not path:
1299 return
1300 # Only emit if the file actually exists
1301 if result.get('exists') is False:
1302 return
1303 page = tool_function_params.get('page')
1305 await event_emitter(
1306 {
1307 'type': f'terminal:{tool_function_name}',
1308 'data': {
1309 'path': path,
1310 **({'page': page} if page else {}),
1311 },
1312 }
1313 )
1314 elif tool_function_name in ('write_file', 'replace_file_content'):
1315 result = parse_terminal_tool_result(tool_result)
1316 path = result.get('path') or tool_function_params.get('path', '')
1317 if not path:
1318 return
1319 await event_emitter(
1320 {
1321 'type': f'terminal:{tool_function_name}',
1322 'data': {'path': path},
1323 }
1324 )
1325 elif tool_function_name == 'run_command':
1326 await event_emitter(
1327 {
1328 'type': 'terminal:run_command',
1329 'data': {},
1330 }
1331 )
1334async def chat_completion_tools_handler(
1335 request: Request, body: dict, extra_params: dict, user: UserModel, models, tools
1336) -> tuple[dict, dict]:
1337 async def get_content_from_response(response) -> Optional[str]:
1338 content = None
1339 if hasattr(response, 'body_iterator'):
1340 async for chunk in response.body_iterator:
1341 data = JSONCodec.loads(chunk.decode('utf-8', 'replace'))
1342 content = data['choices'][0]['message']['content']
1344 # Cleanup any remaining background tasks if necessary
1345 if response.background is not None:
1346 await response.background()
1347 else:
1348 content = response['choices'][0]['message']['content']
1349 return content
1351 def get_tools_function_calling_payload(messages, task_model_id, content):
1352 user_message = get_last_user_message(messages)
1354 if user_message and messages and messages[-1]['role'] == 'user':
1355 # Remove the last user message to avoid duplication
1356 messages = messages[:-1]
1358 recent_messages = messages[-4:] if len(messages) > 4 else messages
1359 chat_history = '\n'.join(
1360 f'{message["role"].upper()}: """{get_content_from_message(message)}"""' for message in recent_messages
1361 )
1363 prompt = f'History:\n{chat_history}\nQuery: {user_message}' if chat_history else f'Query: {user_message}'
1365 return {
1366 'model': task_model_id,
1367 'messages': [
1368 {'role': 'system', 'content': content},
1369 {'role': 'user', 'content': prompt},
1370 ],
1371 'stream': False,
1372 'metadata': {'task': str(TASKS.FUNCTION_CALLING)},
1373 }
1375 event_caller = extra_params['__event_call__']
1376 event_emitter = extra_params['__event_emitter__']
1377 metadata = extra_params['__metadata__']
1379 # One batched SELECT instead of four sequential round trips.
1380 task_config = await Config.get_many(
1381 'task.model.default',
1382 'task.model.external',
1383 'task.tools.prompt_template',
1384 )
1385 task_model_id = get_task_model_id(
1386 body['model'],
1387 task_config.get('task.model.default'),
1388 task_config.get('task.model.external'),
1389 models,
1390 )
1392 skip_files = False
1393 sources = []
1395 specs = [tool['spec'] for tool in tools.values()]
1396 tools_specs = JSONCodec.dumps(specs, ensure_ascii=False)
1398 tools_prompt_template = task_config.get('task.tools.prompt_template')
1399 if tools_prompt_template != '':
1400 template = tools_prompt_template
1401 else:
1402 template = DEFAULT_TOOLS_FUNCTION_CALLING_PROMPT_TEMPLATE
1404 tools_function_calling_prompt = tools_function_calling_generation_template(template, tools_specs)
1405 payload = get_tools_function_calling_payload(body['messages'], task_model_id, tools_function_calling_prompt)
1407 try:
1408 response = await generate_chat_completion(request, form_data=payload, user=user)
1409 log.debug('response=%r', response)
1410 content = await get_content_from_response(response)
1411 log.debug('content=%r', content)
1413 if not content:
1414 return body, {}
1416 try:
1417 content = content[content.find('{') : content.rfind('}') + 1]
1418 if not content:
1419 raise Exception('No JSON object found in the response')
1421 result = JSONCodec.loads(content)
1423 async def tool_call_handler(tool_call):
1424 nonlocal skip_files
1426 log.debug('tool_call=%r', tool_call)
1428 tool_function_name = tool_call.get('name', None)
1429 if tool_function_name not in tools:
1430 log.warning(f'Tool "{tool_function_name}" not found')
1431 return
1433 tool_function_params = tool_call.get('parameters', {})
1435 tool = None
1436 tool_type = ''
1437 direct_tool = False
1439 try:
1440 tool = tools[tool_function_name]
1441 tool_type = tool.get('type', '')
1442 direct_tool = tool.get('direct', False)
1444 spec = tool.get('spec', {})
1445 allowed_params = spec.get('parameters', {}).get('properties', {}).keys()
1446 tool_function_params = {k: v for k, v in tool_function_params.items() if k in allowed_params}
1448 if tool.get('direct', False):
1449 tool_result = await event_caller(
1450 {
1451 'type': 'execute:tool',
1452 'data': {
1453 'id': str(uuid4()),
1454 'name': tool_function_name,
1455 'params': tool_function_params,
1456 'server': tool.get('server', {}),
1457 'session_id': metadata.get('session_id', None),
1458 },
1459 }
1460 )
1461 else:
1462 tool_function = tool['callable']
1463 tool_result = await tool_function(**tool_function_params)
1465 except Exception as e:
1466 tool_result = {'error': str(e)}
1468 tool_result, tool_result_files, tool_result_embeds = await process_tool_result(
1469 request,
1470 tool_function_name,
1471 tool_result,
1472 tool_type,
1473 direct_tool,
1474 metadata,
1475 user,
1476 )
1478 if event_emitter:
1479 await terminal_event_handler(
1480 tool_function_name,
1481 tool_function_params,
1482 tool_result,
1483 event_emitter,
1484 )
1486 if tool_result_files:
1487 for file_item in tool_result_files:
1488 if file_item.get('type') == 'image':
1489 file_item['url'] = await store_tool_result_image(
1490 request, file_item.get('url'), metadata, user
1491 )
1493 await event_emitter(
1494 {
1495 'type': 'files',
1496 'data': {
1497 'files': tool_result_files,
1498 },
1499 }
1500 )
1502 if tool_result_embeds:
1503 await event_emitter(
1504 {
1505 'type': 'embeds',
1506 'data': {
1507 'embeds': tool_result_embeds,
1508 },
1509 }
1510 )
1512 if tool_result:
1513 tool = tools[tool_function_name]
1514 tool_id = tool.get('tool_id', '')
1516 tool_name = f'{tool_id}/{tool_function_name}' if tool_id else f'{tool_function_name}'
1518 # Citation is enabled for this tool
1519 sources.append(
1520 {
1521 'source': {
1522 'name': (f'{tool_name}'),
1523 },
1524 'document': [str(tool_result)],
1525 'metadata': [
1526 {
1527 'source': (f'{tool_name}'),
1528 'parameters': tool_function_params,
1529 }
1530 ],
1531 'tool_result': True,
1532 }
1533 )
1535 if tools[tool_function_name].get('metadata', {}).get('file_handler', False):
1536 skip_files = True
1538 # check if "tool_calls" in result
1539 if result.get('tool_calls'):
1540 for tool_call in result.get('tool_calls'):
1541 await tool_call_handler(tool_call)
1542 else:
1543 await tool_call_handler(result)
1545 except Exception as e:
1546 log.debug('Error: %s', e)
1547 content = None
1548 except Exception as e:
1549 log.debug('Error: %s', e)
1550 content = None
1552 log.debug('tool_contexts: %s', sources)
1554 if skip_files and 'files' in body.get('metadata', {}):
1555 del body['metadata']['files']
1557 return body, {'sources': sources}
1560async def chat_web_search_handler(request: Request, form_data: dict, extra_params: dict, user):
1561 event_emitter = extra_params['__event_emitter__']
1562 await event_emitter(
1563 {
1564 'type': 'status',
1565 'data': {
1566 'action': 'web_search',
1567 'description': 'Searching the web',
1568 'done': False,
1569 },
1570 }
1571 )
1573 messages = form_data['messages']
1574 user_message = get_last_user_message(messages)
1576 queries = []
1577 try:
1578 res = await generate_queries(
1579 request,
1580 {
1581 'model': form_data['model'],
1582 'messages': messages,
1583 'prompt': user_message,
1584 'type': 'web_search',
1585 'chat_id': extra_params.get('__chat_id__'),
1586 },
1587 user,
1588 )
1590 # generate_queries returns a JSONResponse on error (e.g. model not
1591 # found, chat completion failure). Extract the error detail and
1592 # re-raise so the outer except block falls back to using the raw
1593 # user message as the search query.
1594 if isinstance(res, JSONResponse):
1595 try:
1596 error_body = JSONCodec.loads(res.body)
1597 detail = error_body.get('detail', 'Query generation failed')
1598 except Exception:
1599 detail = 'Query generation failed'
1600 raise Exception(detail)
1602 response = res['choices'][0]['message']['content']
1604 try:
1605 bracket_start = response.rfind('{')
1606 bracket_end = response.rfind('}') + 1
1608 if bracket_start == -1 or bracket_end == -1:
1609 raise Exception('No JSON object found in the response')
1611 response = response[bracket_start:bracket_end]
1612 queries = JSONCodec.loads(response)
1613 queries = queries.get('queries', [])
1614 except Exception as e:
1615 queries = [response]
1617 if ENABLE_QUERIES_CACHE:
1618 request.state.cached_queries = queries
1620 except Exception as e:
1621 log.exception(e)
1622 queries = [user_message or '']
1624 # Check if generated queries are empty
1625 if len(queries) == 1 and queries[0].strip() == '':
1626 queries = [user_message or '']
1628 # Check if queries are not found
1629 if len(queries) == 0:
1630 await event_emitter(
1631 {
1632 'type': 'status',
1633 'data': {
1634 'action': 'web_search',
1635 'description': 'No search query generated',
1636 'done': True,
1637 },
1638 }
1639 )
1640 return form_data
1642 await event_emitter(
1643 {
1644 'type': 'status',
1645 'data': {
1646 'action': 'web_search_queries_generated',
1647 'queries': queries,
1648 'done': False,
1649 },
1650 }
1651 )
1653 try:
1654 results = await process_web_search(
1655 request,
1656 SearchForm(queries=queries),
1657 user=user,
1658 )
1660 if results:
1661 files = form_data.get('files', [])
1663 if results.get('collection_names'):
1664 for col_idx, collection_name in enumerate(results.get('collection_names')):
1665 files.append(
1666 {
1667 'collection_name': collection_name,
1668 'name': ', '.join(queries),
1669 'type': 'web_search',
1670 'urls': results['filenames'],
1671 'queries': queries,
1672 }
1673 )
1674 elif results.get('docs'):
1675 # Invoked when bypass embedding and retrieval is set to True
1676 docs = results['docs']
1677 files.append(
1678 {
1679 'docs': docs,
1680 'name': ', '.join(queries),
1681 'type': 'web_search',
1682 'urls': results['filenames'],
1683 'queries': queries,
1684 }
1685 )
1687 form_data['files'] = files
1689 await event_emitter(
1690 {
1691 'type': 'status',
1692 'data': {
1693 'action': 'web_search',
1694 'description': 'Searched {{count}} sites',
1695 'urls': results['filenames'],
1696 'items': results.get('items', []),
1697 'done': True,
1698 },
1699 }
1700 )
1701 else:
1702 await event_emitter(
1703 {
1704 'type': 'status',
1705 'data': {
1706 'action': 'web_search',
1707 'description': 'No search results found',
1708 'done': True,
1709 'error': True,
1710 },
1711 }
1712 )
1714 except Exception as e:
1715 log.exception(e)
1716 detail = e.detail if isinstance(e, HTTPException) else None
1717 await event_emitter(
1718 {
1719 'type': 'status',
1720 'data': {
1721 'action': 'web_search',
1722 'description': (str(detail) if detail else 'An error occurred while searching the web'),
1723 'queries': queries,
1724 'done': True,
1725 'error': True,
1726 },
1727 }
1728 )
1730 return form_data
1733def get_images_from_messages(message_list):
1734 images = []
1736 for message in reversed(message_list):
1737 message_images = []
1738 for file in message.get('files', []):
1739 if file.get('type') == 'image':
1740 message_images.append(file.get('url'))
1741 elif is_raster_image_content_type(file.get('content_type')):
1742 message_images.append(file.get('url'))
1744 if message_images:
1745 images.append(message_images)
1747 return images
1750async def get_image_urls(delta_images, request, metadata, user) -> list[str]:
1751 if not isinstance(delta_images, list):
1752 return []
1754 image_urls = []
1755 for img in delta_images:
1756 if not isinstance(img, dict) or img.get('type') != 'image_url':
1757 continue
1759 url = img.get('image_url', {}).get('url')
1760 if not url:
1761 continue
1763 if url.startswith('data:image/png;base64'):
1764 url = await get_image_url_from_base64(request, url, metadata, user)
1766 image_urls.append(url)
1768 return image_urls
1771async def add_file_context(messages: list, chat_id: str, user) -> list:
1772 """
1773 Add file URLs to messages for native function calling.
1774 """
1775 if not is_saved_chat_id(chat_id):
1776 return messages
1778 chat = await Chats.get_chat_by_id_and_user_id(chat_id, user.id)
1779 if not chat:
1780 return messages
1782 history = chat.chat.get('history', {})
1783 stored_messages = get_message_list(history.get('messages', {}), history.get('currentId'))
1785 def format_file_tag(file):
1786 # Every file reaching here has a url or a chat id, so id is always set.
1787 attrs = f'type="{file.get("type", "file")}" id="{file.get("id") or file.get("url")}"'
1788 if file.get('url'):
1789 attrs += f' url="{file["url"]}"'
1790 if file.get('content_type'):
1791 attrs += f' content_type="{file["content_type"]}"'
1792 if file.get('name'):
1793 attrs += f' name="{file["name"]}"'
1794 return f'<file {attrs}/>'
1796 # Pair only user-role messages from both lists to avoid misalignment.
1797 # After process_messages_with_output(), assistant messages with tool calls
1798 # are expanded into multiple messages (assistant + tool results), making
1799 # the payload message list longer than the stored message list. A naive
1800 # positional zip() would pair user messages with wrong stored messages,
1801 # causing later images to lose their file context (see #21878).
1802 user_messages = [m for m in messages if m.get('role') == 'user']
1803 stored_user_messages = [m for m in stored_messages if m.get('role') == 'user']
1805 for message, stored_message in zip(user_messages, stored_user_messages):
1806 # Chat references carry no url - they are addressed by id via view_chat.
1807 attached_files = [
1808 file
1809 for file in stored_message.get('files', [])
1810 if (file.get('url') and not file.get('url').startswith('data:'))
1811 or (file.get('type') == 'chat' and file.get('id'))
1812 ]
1813 if not attached_files:
1814 continue
1816 file_tags = [format_file_tag(file) for file in attached_files]
1817 file_context = '<attached_files>\n' + '\n'.join(file_tags) + '\n</attached_files>\n\n'
1819 content = message.get('content', '')
1820 if isinstance(content, list):
1821 message['content'] = [{'type': 'text', 'text': file_context}] + content
1822 else:
1823 message['content'] = file_context + content
1825 return messages
1828async def chat_image_generation_handler(request: Request, form_data: dict, extra_params: dict, user):
1829 metadata = extra_params.get('__metadata__', {})
1830 chat_id = metadata.get('chat_id', None)
1831 __event_emitter__ = extra_params.get('__event_emitter__', None)
1833 if not chat_id or not isinstance(chat_id, str) or not __event_emitter__:
1834 return form_data
1836 is_channel_chat = chat_id.startswith('channel:')
1837 image_metadata = {
1838 'message_id': metadata.get('message_id', None),
1839 **({'channel_id': chat_id.removeprefix('channel:')} if is_channel_chat else {'chat_id': chat_id}),
1840 }
1842 if not is_saved_chat_id(chat_id):
1843 message_list = form_data.get('messages', [])
1844 else:
1845 chat = await Chats.get_chat_by_id_and_user_id(chat_id, user.id)
1847 messages_map = chat.chat.get('history', {}).get('messages', {})
1848 message_id = chat.chat.get('history', {}).get('currentId')
1849 message_list = get_message_list(messages_map, message_id)
1851 user_message = get_last_user_message(message_list)
1853 prompt = user_message
1854 message_images = get_images_from_messages(message_list)
1856 # Limit to first 2 sets of images
1857 # We may want to change this in the future to allow more images
1858 input_images = []
1859 for idx, images in enumerate(message_images):
1860 if idx >= 2:
1861 break
1862 for image in images:
1863 input_images.append(image)
1865 # Called directly, bypassing the /images routes that enforce these switches.
1866 editing = len(input_images) > 0 and await Config.get('images.edit.enable')
1867 if not editing and not await Config.get('image_generation.enable'):
1868 return form_data
1870 if is_saved_chat_id(chat_id):
1871 await __event_emitter__(
1872 {
1873 'type': 'status',
1874 'data': {'description': 'Creating image', 'done': False},
1875 }
1876 )
1878 system_message_content = ''
1880 if editing:
1881 # Edit image(s)
1882 try:
1883 images = await image_edits(
1884 request=request,
1885 form_data=EditImageForm(**{'prompt': prompt, 'image': input_images}),
1886 metadata=image_metadata,
1887 user=user,
1888 )
1890 await __event_emitter__(
1891 {
1892 'type': 'status',
1893 'data': {'description': 'Image created', 'done': True},
1894 }
1895 )
1897 await __event_emitter__(
1898 {
1899 'type': 'files',
1900 'data': {
1901 'files': [
1902 {
1903 'type': 'image',
1904 **image,
1905 }
1906 for image in images
1907 ]
1908 },
1909 }
1910 )
1912 system_message_content = '<context>The requested image has been edited and created and is now being shown to the user. Let them know that it has been generated.</context>'
1913 except Exception as e:
1914 log.debug(e)
1916 error_message = ''
1917 if isinstance(e, HTTPException):
1918 if e.detail and isinstance(e.detail, dict):
1919 error_message = e.detail.get('message', str(e.detail))
1920 else:
1921 error_message = str(e.detail)
1923 await __event_emitter__(
1924 {
1925 'type': 'status',
1926 'data': {
1927 'description': f'An error occurred while generating an image',
1928 'done': True,
1929 },
1930 }
1931 )
1933 system_message_content = f'<context>Image generation was attempted but failed. The system is currently unable to generate the image. Tell the user that the following error occurred: {error_message}</context>'
1935 elif not await Config.get('image_generation.enable'):
1936 await __event_emitter__(
1937 {
1938 'type': 'status',
1939 'data': {
1940 'description': 'Image generation is disabled',
1941 'done': True,
1942 },
1943 }
1944 )
1946 system_message_content = '<context>Image generation was requested but the feature is currently disabled by the administrator, so no image was created. Let the user know that image generation is currently unavailable.</context>'
1948 else:
1949 # Create image(s)
1950 if await Config.get('image_generation.prompt.enable'):
1951 try:
1952 res = await generate_image_prompt(
1953 request,
1954 {
1955 'model': form_data['model'],
1956 'messages': form_data['messages'],
1957 'chat_id': metadata.get('chat_id'),
1958 },
1959 user,
1960 )
1962 # Handle JSONResponse from error paths
1963 if isinstance(res, JSONResponse):
1964 try:
1965 error_body = JSONCodec.loads(res.body)
1966 detail = error_body.get('detail', 'Image prompt generation failed')
1967 except Exception:
1968 detail = 'Image prompt generation failed'
1969 raise Exception(detail)
1971 response = res['choices'][0]['message']['content']
1973 try:
1974 bracket_start = response.rfind('{')
1975 bracket_end = response.rfind('}') + 1
1977 if bracket_start == -1 or bracket_end == -1:
1978 raise Exception('No JSON object found in the response')
1980 response = response[bracket_start:bracket_end]
1981 response = JSONCodec.loads(response)
1982 prompt = response.get('prompt', [])
1983 except Exception as e:
1984 prompt = user_message
1986 except Exception as e:
1987 log.exception(e)
1988 prompt = user_message
1990 try:
1991 images = await image_generations(
1992 request=request,
1993 form_data=CreateImageForm(**{'prompt': prompt}),
1994 metadata=image_metadata,
1995 user=user,
1996 )
1998 await __event_emitter__(
1999 {
2000 'type': 'status',
2001 'data': {'description': 'Image created', 'done': True},
2002 }
2003 )
2005 await __event_emitter__(
2006 {
2007 'type': 'files',
2008 'data': {
2009 'files': [
2010 {
2011 'type': 'image',
2012 **image,
2013 }
2014 for image in images
2015 ]
2016 },
2017 }
2018 )
2020 system_message_content = '<context>The requested image has been created by the system successfully and is now being shown to the user. Let the user know that the image they requested has been generated and is now shown in the chat.</context>'
2021 except Exception as e:
2022 log.debug(e)
2024 error_message = ''
2025 if isinstance(e, HTTPException):
2026 if e.detail and isinstance(e.detail, dict):
2027 error_message = e.detail.get('message', str(e.detail))
2028 else:
2029 error_message = str(e.detail)
2031 await __event_emitter__(
2032 {
2033 'type': 'status',
2034 'data': {
2035 'description': f'An error occurred while generating an image',
2036 'done': True,
2037 },
2038 }
2039 )
2041 system_message_content = f'<context>Image generation was attempted but failed because of an error. The system is currently unable to generate the image. Tell the user that the following error occurred: {error_message}</context>'
2043 if system_message_content:
2044 form_data['messages'] = add_or_update_system_message(system_message_content, form_data['messages'])
2046 return form_data
2049async def chat_completion_files_handler(
2050 request: Request, body: dict, extra_params: dict, user: UserModel
2051) -> tuple[dict, dict[str, list]]:
2052 __event_emitter__ = extra_params['__event_emitter__']
2053 sources = []
2055 files = [item for item in (body.get('metadata', {}).get('files', None) or []) if item.get('type') != 'filesystem']
2056 if files:
2057 # Check if all files are in full context mode
2058 all_full_context = all(item.get('context') == 'full' for item in files)
2060 queries = []
2061 if not all_full_context:
2062 try:
2063 queries_response = await generate_queries(
2064 request,
2065 {
2066 'model': body['model'],
2067 'messages': body['messages'],
2068 'type': 'retrieval',
2069 'chat_id': body.get('metadata', {}).get('chat_id'),
2070 },
2071 user,
2072 )
2073 queries_response = queries_response['choices'][0]['message']['content']
2075 try:
2076 bracket_start = queries_response.rfind('{')
2077 bracket_end = queries_response.rfind('}') + 1
2079 if bracket_start == -1 or bracket_end == -1:
2080 raise Exception('No JSON object found in the response')
2082 queries_response = queries_response[bracket_start:bracket_end]
2083 queries_response = JSONCodec.loads(queries_response)
2084 except Exception as e:
2085 queries_response = {'queries': [queries_response]}
2087 queries = queries_response.get('queries', [])
2088 except Exception:
2089 pass
2091 await __event_emitter__(
2092 {
2093 'type': 'status',
2094 'data': {
2095 'action': 'queries_generated',
2096 'queries': queries,
2097 'done': False,
2098 },
2099 }
2100 )
2102 if len(queries) == 0:
2103 queries = [get_last_user_message(body['messages']) or '']
2105 try:
2106 # One batched SELECT instead of six sequential round trips.
2107 rag_config = await Config.get_many(
2108 'rag.top_k',
2109 'rag.top_k_reranker',
2110 'rag.relevance_threshold',
2111 'rag.hybrid_bm25_weight',
2112 'rag.enable_hybrid_search',
2113 'rag.full_context',
2114 )
2115 # Directly await async get_sources_from_items (no thread needed - fully async now)
2116 sources = await get_sources_from_items(
2117 request=request,
2118 items=files,
2119 queries=queries,
2120 embedding_function=lambda query, prefix: request.app.state.EMBEDDING_FUNCTION(
2121 query, prefix=prefix, user=user
2122 ),
2123 k=rag_config.get('rag.top_k'),
2124 reranking_function=(
2125 (lambda query, documents: request.app.state.RERANKING_FUNCTION(query, documents, user=user))
2126 if request.app.state.RERANKING_FUNCTION
2127 else None
2128 ),
2129 k_reranker=rag_config.get('rag.top_k_reranker'),
2130 r=rag_config.get('rag.relevance_threshold'),
2131 hybrid_bm25_weight=rag_config.get('rag.hybrid_bm25_weight'),
2132 hybrid_search=rag_config.get('rag.enable_hybrid_search'),
2133 full_context=all_full_context or rag_config.get('rag.full_context'),
2134 user=user,
2135 )
2136 except Exception as e:
2137 log.exception(e)
2139 log.debug('rag_contexts:sources: %s', sources)
2141 unique_ids = set()
2142 for source in sources or []:
2143 if not source or len(source.keys()) == 0:
2144 continue
2146 documents = source.get('document') or []
2147 metadatas = source.get('metadata') or []
2148 src_info = source.get('source') or {}
2150 for index, _ in enumerate(documents):
2151 metadata = metadatas[index] if index < len(metadatas) else None
2152 _id = (metadata or {}).get('source') or (src_info or {}).get('id') or 'N/A'
2153 unique_ids.add(_id)
2155 sources_count = len(unique_ids)
2156 await __event_emitter__(
2157 {
2158 'type': 'status',
2159 'data': {
2160 'action': 'sources_retrieved',
2161 'count': sources_count,
2162 'done': True,
2163 },
2164 }
2165 )
2167 return body, {'sources': sources}
2170async def convert_url_images_to_base64(form_data, user=None):
2171 messages = form_data.get('messages', [])
2173 for message in messages:
2174 content = message.get('content')
2175 if not isinstance(content, list):
2176 continue
2178 new_content = []
2180 for item in content:
2181 if not isinstance(item, dict) or item.get('type') not in ('image_url', 'input_image'):
2182 new_content.append(item)
2183 continue
2185 image_url_data = item.get('image_url', {})
2186 if isinstance(image_url_data, dict):
2187 image_url = image_url_data.get('url') or ''
2188 elif isinstance(image_url_data, str):
2189 image_url = image_url_data
2190 else:
2191 image_url = ''
2192 if image_url.startswith('data:image/'):
2193 new_content.append(item)
2194 continue
2196 try:
2197 base64_data = await get_image_base64_from_url(image_url, user=user)
2198 if base64_data and isinstance(image_url_data, str):
2199 new_content.append({**item, 'image_url': base64_data})
2200 elif base64_data:
2201 image_url_payload = {'url': base64_data}
2202 if isinstance(image_url_data, dict) and image_url_data.get('detail'):
2203 image_url_payload['detail'] = image_url_data['detail']
2204 new_content.append(
2205 {
2206 'type': item['type'],
2207 'image_url': image_url_payload,
2208 }
2209 )
2210 else:
2211 new_content.append(item)
2212 except Exception as e:
2213 log.debug('Error converting image URL to base64: %s', e)
2214 new_content.append(item)
2216 message['content'] = new_content
2218 return form_data
2221MESSAGE_REPLAY_KEYS = ('id', 'role', 'content', 'output', 'files', 'contextSummary', 'usage', 'model')
2224async def load_messages_from_db(chat_id: str, message_id: str) -> Optional[list[dict]]:
2225 """
2226 Load the message chain from DB up to message_id,
2227 keeping only fields needed to rebuild the LLM payload.
2228 """
2229 messages_map = await Chats.get_messages_map_by_chat_id(chat_id)
2230 if not messages_map:
2231 return None
2233 db_messages = get_message_list(messages_map, message_id)
2234 if not db_messages:
2235 return None
2237 return [
2238 {k: v for k, v in msg.items() if k in MESSAGE_REPLAY_KEYS}
2239 for msg in db_messages
2240 if not (
2241 msg.get('role') == 'assistant' and msg.get('error') and not msg.get('content') and not msg.get('output')
2242 )
2243 ]
2246def get_reasoning_format(model: dict) -> str | None:
2247 """
2248 Determine how reasoning should be included in reconstructed messages.
2250 Returns:
2251 'thinking': Ollama expects reasoning in the native thinking field.
2252 'think_tags': wrap reasoning in <think> tags inside content.
2253 'reasoning_content': llama.cpp supports reasoning_content as a top-level field.
2254 None: skip reasoning (safe default for strict providers).
2255 """
2256 provider = model.get('provider', '')
2257 if model.get('owned_by') == 'ollama':
2258 return 'thinking'
2259 if provider == 'llama.cpp':
2260 return 'reasoning_content'
2261 return None
2264def strip_reasoning_details(output: list) -> list:
2265 return [
2266 {key: value for key, value in item.items() if key != 'reasoning_details'} if isinstance(item, dict) else item
2267 for item in output
2268 ]
2271def process_messages_with_output(
2272 messages: list[dict],
2273 reasoning_format: str | None = None,
2274) -> list[dict]:
2275 """
2276 Process messages with OR-aligned output items for LLM consumption.
2278 For assistant messages with 'output' field, produces properly formatted
2279 OpenAI-style messages (tool_calls + tool results). Strips 'output' before LLM.
2280 """
2281 processed = []
2283 for message in messages:
2284 if message.get('role') == 'assistant' and message.get('output'):
2285 # Use output items for clean OpenAI-format messages
2286 output_messages = convert_output_to_messages(
2287 message['output'],
2288 raw=True,
2289 reasoning_format=reasoning_format,
2290 flatten_tool_images=True,
2291 )
2292 if output_messages:
2293 processed.extend(output_messages)
2294 continue
2295 if not message.get('content'):
2296 continue
2298 clean_message = dict(message)
2299 for key in ('id', 'files', 'output', 'model', 'contextSummary', 'context_summary', 'usage'):
2300 clean_message.pop(key, None)
2301 processed.append(clean_message)
2303 return processed
2306def sanitize_tool_pairs(messages: list[dict]) -> list[dict]:
2307 tool_result_ids = {
2308 message.get('tool_call_id')
2309 for message in messages
2310 if message.get('role') == 'tool' and message.get('tool_call_id')
2311 }
2313 tool_call_ids = {
2314 tool_call.get('id')
2315 for message in messages
2316 for tool_call in (message.get('tool_calls') or [])
2317 if message.get('role') == 'assistant' and tool_call.get('id')
2318 }
2320 sanitized = []
2321 for message in messages:
2322 if message.get('role') == 'assistant' and message.get('tool_calls'):
2323 kept = [
2324 tool_call for tool_call in message.get('tool_calls') or [] if tool_call.get('id') in tool_result_ids
2325 ]
2326 if kept:
2327 sanitized.append({**message, 'tool_calls': kept})
2328 else:
2329 clean = dict(message)
2330 clean.pop('tool_calls', None)
2331 clean.pop('reasoning_items', None)
2332 if clean.get('content'):
2333 sanitized.append(clean)
2334 elif message.get('role') != 'tool' or message.get('tool_call_id') in tool_call_ids:
2335 sanitized.append(message)
2337 return sanitized
2340async def connect_mcp_server(
2341 request,
2342 server_id: str,
2343 user,
2344 metadata: dict,
2345 extra_params: dict,
2346) -> tuple[MCPClient, list[dict]] | None:
2347 """Resolve an MCP server connection, authenticate, and return (client, tool_specs).
2349 Returns None if the server is not found or access is denied.
2350 """
2351 mcp_server_connection = None
2352 for server_connection in await Config.get('tool_server.connections', []):
2353 if server_connection.get('type', '') == 'mcp' and (server_connection.get('info') or {}).get('id') == server_id:
2354 mcp_server_connection = server_connection
2355 break
2357 if not mcp_server_connection:
2358 log.error(f'MCP server with id {server_id} not found')
2359 return None
2361 if not await has_connection_access(user, mcp_server_connection):
2362 log.warning(f'Access denied to MCP server {server_id} for user {user.id}')
2363 return None
2365 headers, _ = await build_tool_server_headers(
2366 mcp_server_connection,
2367 request,
2368 user,
2369 server_id=server_id,
2370 metadata=metadata,
2371 extra_params=extra_params,
2372 )
2374 client = MCPClient()
2375 await client.connect(
2376 url=mcp_server_connection.get('url', ''),
2377 headers=headers if headers else None,
2378 )
2380 function_name_filter_list = mcp_server_connection.get('config', {}).get('function_name_filter_list', '')
2381 if isinstance(function_name_filter_list, str):
2382 function_name_filter_list = function_name_filter_list.split(',')
2384 tool_specs = await client.list_tool_specs()
2385 if function_name_filter_list:
2386 tool_specs = [spec for spec in tool_specs if is_string_allowed(spec['name'], function_name_filter_list)]
2388 return client, tool_specs
2391async def process_chat_payload(request, form_data, user, metadata, model):
2392 # Ensure chat_id is always a string — external API clients may omit it.
2393 if not isinstance(metadata.get('chat_id'), str):
2394 metadata['chat_id'] = ''
2396 # Pipeline Inlet -> Filter Inlet -> Chat Memory -> Chat Web Search -> Chat Image Generation
2397 # -> Chat Code Interpreter (Form Data Update) -> (Default) Chat Tools Function Calling
2398 # -> Chat Files
2400 # Arena model resolution — pick the sub-model now so all downstream
2401 # processing (knowledge, capabilities, tools, params) uses its settings
2402 # instead of the empty arena wrapper.
2403 if model.get('owned_by') == 'arena':
2404 arena_model_ids = model.get('info', {}).get('meta', {}).get('model_ids')
2405 arena_filter_mode = model.get('info', {}).get('meta', {}).get('filter_mode')
2406 if arena_model_ids and arena_filter_mode == 'exclude':
2407 arena_model_ids = [
2408 available_model['id']
2409 for available_model in request.app.state.MODELS.values()
2410 if available_model.get('owned_by') != 'arena' and available_model['id'] not in arena_model_ids
2411 ]
2413 if isinstance(arena_model_ids, list) and arena_model_ids:
2414 selected_model_id = random.choice(arena_model_ids)
2415 else:
2416 arena_model_ids = [
2417 available_model['id']
2418 for available_model in request.app.state.MODELS.values()
2419 if available_model.get('owned_by') != 'arena'
2420 ]
2421 selected_model_id = random.choice(arena_model_ids)
2423 selected_model = request.app.state.MODELS.get(selected_model_id)
2424 if selected_model:
2425 model = selected_model
2426 form_data['model'] = selected_model_id
2427 metadata['selected_model_id'] = selected_model_id
2429 # Captured before apply_params_to_form_data pops 'params'; populates metadata['system_prompt'] below
2430 model_system_prompt = (form_data.get('params') or {}).get('system')
2432 form_data = apply_params_to_form_data(form_data, model)
2433 log.debug('form_data: %s', form_data)
2435 # Guided regeneration: extract before it reaches the LLM provider
2436 regeneration_prompt = form_data.pop('regeneration_prompt', None)
2438 # Load messages from DB when available — DB preserves structured 'output' items
2439 # which the frontend strips, causing tool calls to be merged into content.
2440 chat_id = metadata.get('chat_id')
2441 user_message_id = metadata.get('user_message_id')
2443 if is_saved_chat_id(chat_id) and user_message_id:
2444 db_messages = await load_messages_from_db(chat_id, user_message_id)
2445 if db_messages:
2446 # Continue: frontend sends assistant_message_id when continuing
2447 # an existing response. Load its content so the LLM sees prior output.
2448 assistant_message_id = metadata.get('assistant_message_id')
2449 if assistant_message_id:
2450 assistant_message = await Chats.get_message_by_id_and_message_id(chat_id, assistant_message_id)
2451 if assistant_message and (assistant_message.get('content') or assistant_message.get('output')):
2452 db_messages.append({k: v for k, v in assistant_message.items() if k in MESSAGE_REPLAY_KEYS})
2454 system_message = get_system_message(form_data.get('messages', []))
2455 form_data['messages'] = [system_message, *db_messages] if system_message else db_messages
2457 # Inject image files into content as image_url parts (mirrors frontend logic)
2458 for message in form_data['messages']:
2459 image_files = [
2460 f
2461 for f in message.get('files', [])
2462 if f.get('type') == 'image' or is_raster_image_content_type(f.get('content_type'))
2463 ]
2464 if message.get('role') == 'user' and image_files:
2465 text_content = message.get('content', '')
2466 if isinstance(text_content, str):
2467 message['content'] = [
2468 {'type': 'text', 'text': text_content},
2469 *[
2470 {
2471 'type': 'image_url',
2472 'image_url': {'url': f['url']},
2473 }
2474 for f in image_files
2475 if f.get('url')
2476 ],
2477 ]
2478 # Strip files field — it's been incorporated into content
2479 message.pop('files', None)
2481 if regeneration_prompt:
2482 form_data['messages'].append({'role': 'user', 'content': regeneration_prompt})
2484 if is_saved_chat_id(chat_id) and user_message_id:
2485 if getattr(request.state, 'direct', False) and hasattr(request.state, 'model'):
2486 compaction_models = {
2487 **dict(request.app.state.MODELS.items()),
2488 request.state.model['id']: request.state.model,
2489 }
2490 else:
2491 compaction_models = request.app.state.MODELS
2493 system_message = get_system_message(form_data.get('messages', []))
2494 system_prompt = get_content_from_message(system_message) if system_message else ''
2496 try:
2497 form_data['messages'], context_summary, _ = await compact_messages_for_request(
2498 request,
2499 user,
2500 form_data.get('messages', []),
2501 metadata,
2502 form_data.get('model'),
2503 compaction_models,
2504 system_prompt,
2505 )
2506 if context_summary:
2507 form_data['messages'] = add_or_update_system_message(
2508 f'[CONVERSATION SUMMARY]\n{context_summary}',
2509 form_data['messages'],
2510 append=True,
2511 )
2512 except Exception:
2513 log.exception('Context compaction failed; continuing with full chat history')
2515 # Process messages with OR-aligned output items for clean LLM messages
2516 for message in form_data.get('messages', []):
2517 output = message.get('output')
2518 # reasoning_details can be model/provider-bound, so only replay them
2519 # for output produced by the same model.
2520 if message.get('role') == 'assistant' and message.get('model') != model['id'] and isinstance(output, list):
2521 message['output'] = strip_reasoning_details(output)
2523 form_data['messages'] = process_messages_with_output(
2524 form_data.get('messages', []),
2525 reasoning_format=get_reasoning_format(model),
2526 )
2527 form_data['messages'] = sanitize_tool_pairs(form_data['messages'])
2529 system_message = get_system_message(form_data.get('messages', []))
2530 if system_message: # Chat Controls/User Settings
2531 try:
2532 form_data = await apply_system_prompt_to_body(
2533 system_message.get('content'), form_data, metadata, user, replace=True
2534 ) # Required to handle system prompt variables
2535 except Exception:
2536 pass
2538 form_data = await convert_url_images_to_base64(form_data, user=user)
2540 event_emitter = await get_event_emitter(metadata)
2541 event_caller = await get_event_call(metadata)
2543 extra_params = {
2544 '__event_emitter__': event_emitter,
2545 '__event_call__': event_caller,
2546 '__user__': user.model_dump() if isinstance(user, UserModel) else {},
2547 '__metadata__': metadata,
2548 '__oauth_token__': await get_system_oauth_token(request, user),
2549 '__request__': request,
2550 '__model__': model,
2551 '__chat_id__': metadata.get('chat_id'),
2552 '__message_id__': metadata.get('message_id'),
2553 }
2554 # Initialize events to store additional event to be sent to the client
2555 # Initialize contexts and citation
2556 if getattr(request.state, 'direct', False) and hasattr(request.state, 'model'):
2557 models = {
2558 request.state.model['id']: request.state.model,
2559 }
2560 else:
2561 models = request.app.state.MODELS
2563 task_model_id = get_task_model_id(
2564 form_data['model'],
2565 await Config.get('task.model.default'),
2566 await Config.get('task.model.external'),
2567 models,
2568 )
2570 events = []
2571 sources = []
2573 # Folder "Project" handling
2574 # Check if the request has chat_id and is inside of a folder
2575 # Uses lightweight column query — only fetches folder_id, not the full chat JSON blob
2576 chat_id = metadata.get('chat_id', None)
2577 folder_id = None
2578 if user and is_saved_chat_id(chat_id):
2579 folder_id = await Chats.get_chat_folder_id(chat_id, user.id)
2581 # Fallback: use folder_id from metadata (temporary chats have no DB record)
2582 if not folder_id:
2583 folder_id = metadata.get('folder_id', None)
2585 if folder_id and user:
2586 folder = await Folders.get_folder_by_id(folder_id)
2587 if folder and user.role != 'admin' and not await has_folder_access(user.id, folder, 'read', db=None):
2588 folder = None
2590 if folder and folder.data:
2591 if 'system_prompt' in folder.data:
2592 form_data = await apply_system_prompt_to_body(folder.data['system_prompt'], form_data, metadata, user)
2593 if 'files' in folder.data:
2594 if metadata.get('params', {}).get('function_calling') == 'legacy':
2595 form_data['files'] = [
2596 {'type': 'folder', 'id': folder.id},
2597 *form_data.get('files', []),
2598 ]
2599 else:
2600 # Native FC: skip RAG injection, builtin tools
2601 # will read folder knowledge from metadata.
2602 metadata['folder_knowledge'] = await get_owner_accessible_folder_files(folder)
2604 # Model "Knowledge" handling
2605 user_message = get_last_user_message(form_data['messages'])
2606 model_knowledge = model.get('info', {}).get('meta', {}).get('knowledge', False)
2608 if model_knowledge and metadata.get('params', {}).get('function_calling') == 'legacy':
2609 await event_emitter(
2610 {
2611 'type': 'status',
2612 'data': {
2613 'action': 'knowledge_search',
2614 'query': user_message,
2615 'done': False,
2616 },
2617 }
2618 )
2620 knowledge_files = []
2621 for item in model_knowledge:
2622 if item.get('collection_name'):
2623 knowledge_files.append(
2624 {
2625 'id': item.get('collection_name'),
2626 'name': item.get('name'),
2627 'legacy': True,
2628 }
2629 )
2630 elif item.get('collection_names'):
2631 knowledge_files.append(
2632 {
2633 'name': item.get('name'),
2634 'type': 'collection',
2635 'collection_names': item.get('collection_names'),
2636 'legacy': True,
2637 }
2638 )
2639 else:
2640 knowledge_files.append(item)
2642 files = form_data.get('files', [])
2643 files.extend(knowledge_files)
2644 form_data['files'] = files
2646 variables = form_data.pop('variables', None)
2647 payload_tools = form_data.get('tools', None) # snapshot before filters
2649 # Process the form_data through the pipeline
2650 try:
2651 form_data = await process_pipeline_inlet_filter(request, form_data, user, models)
2652 except Exception as e:
2653 raise e
2655 filter_functions = []
2656 filter_context = get_filter_context(request) if ENABLE_PLUGINS else None
2657 if ENABLE_PLUGINS:
2658 try:
2659 filter_functions = await get_filter_functions(request, model, metadata.get('filter_ids', []))
2661 form_data, flags = await process_filter_functions(
2662 request=request,
2663 filter_context=filter_context,
2664 filter_functions=filter_functions,
2665 filter_type='inlet',
2666 form_data=form_data,
2667 extra_params=extra_params,
2668 )
2669 except Exception as e:
2670 raise Exception(f'{e}')
2672 features = form_data.pop('features', None) or {}
2673 extra_params['__features__'] = features
2674 if features:
2675 if 'voice' in features and features['voice']:
2676 if await Config.get('task.voice.prompt.enable'):
2677 template = await Config.get('task.voice.prompt_template')
2678 if not template:
2679 template = DEFAULT_VOICE_MODE_PROMPT_TEMPLATE
2681 form_data['messages'] = add_or_update_system_message(
2682 template,
2683 form_data['messages'],
2684 )
2686 if (
2687 'memory' in features
2688 and features['memory']
2689 and await Config.get('memories.enable')
2690 and await Config.get('memories.system_context.enable')
2691 ):
2692 # features is client-supplied; re-check the permission the native FC path enforces.
2693 if getattr(user, 'role', None) == 'admin' or await has_permission(
2694 getattr(user, 'id', ''),
2695 'features.memories',
2696 await Config.get('user.permissions'),
2697 ):
2698 form_data = await add_memory_context(request, form_data, user, model)
2700 if 'web_search' in features and features['web_search'] and await Config.get('web.search.enable'):
2701 # features is client-supplied; re-check the permission the native FC path enforces.
2702 if getattr(user, 'role', None) == 'admin' or await has_permission(
2703 getattr(user, 'id', ''),
2704 'features.web_search',
2705 await Config.get('user.permissions'),
2706 ):
2707 # Skip forced RAG web search when native FC is enabled - model can use web_search tool
2708 if metadata.get('params', {}).get('function_calling') == 'legacy':
2709 form_data = await chat_web_search_handler(request, form_data, extra_params, user)
2711 if 'image_generation' in features and features['image_generation']:
2712 # features is client-supplied; re-check the permission the direct /images routes enforce.
2713 if getattr(user, 'role', None) == 'admin' or await has_permission(
2714 getattr(user, 'id', ''),
2715 'features.image_generation',
2716 await Config.get('user.permissions'),
2717 ):
2718 # Skip forced image generation when native FC is enabled - model can use generate_image tool
2719 if metadata.get('params', {}).get('function_calling') == 'legacy':
2720 form_data = await chat_image_generation_handler(request, form_data, extra_params, user)
2722 if 'code_interpreter' in features and features['code_interpreter']:
2723 engine = await Config.get('code_interpreter.engine', 'pyodide')
2725 # Skip XML-tag prompt injection when native FC is enabled —
2726 # execute_code will be injected as a builtin tool instead
2727 if metadata.get('params', {}).get('function_calling') == 'legacy':
2728 ci_prompt_template = await Config.get('code_interpreter.prompt_template')
2729 prompt = ci_prompt_template if ci_prompt_template != '' else DEFAULT_CODE_INTERPRETER_PROMPT
2731 # Append filesystem awareness only for pyodide engine
2732 if engine != 'jupyter':
2733 prompt += CODE_INTERPRETER_PYODIDE_PROMPT
2735 form_data['messages'] = add_or_update_user_message(
2736 prompt,
2737 form_data['messages'],
2738 )
2739 else:
2740 # Native FC: tool docstring can't be dynamic, so inject
2741 # filesystem context into the system message for pyodide
2742 # engine. Appending to the system prompt (instead of the
2743 # user message) keeps it in the stable cached prefix so
2744 # providers with prefix caching don't re-bill the full
2745 # conversation on every turn.
2746 if engine != 'jupyter':
2747 form_data['messages'] = add_or_update_system_message(
2748 CODE_INTERPRETER_PYODIDE_PROMPT,
2749 form_data['messages'],
2750 append=True,
2751 )
2753 tool_ids = form_data.pop('tool_ids', None)
2754 terminal_id = form_data.pop('terminal_id', None)
2755 files = form_data.pop('files', None)
2756 form_data.pop('folder_id', None)
2757 metadata['terminal_id'] = terminal_id
2758 skill_authoring_allowed = bool(terminal_id) and has_prior_real_chat_content(form_data.get('messages', []))
2759 skill_create_denial_reason = 'empty_chat' if terminal_id else 'disabled'
2760 apply_skills_create_prompt(
2761 form_data.get('messages', []),
2762 allowed=skill_authoring_allowed,
2763 denial_reason=skill_create_denial_reason,
2764 )
2766 # If the original caller provided tools, use them as-is (skip resolution).
2767 # Otherwise, save any tools that filter inlets added for merging later.
2768 inlet_filter_tools = None if payload_tools is not None else form_data.get('tools', None)
2770 # Mentioned skills get full content; selected/default skills can be loaded through view_skill.
2771 mentioned_skill_ids = extract_skill_ids_from_messages(form_data.get('messages', []))
2772 skill_ids = sorted(
2773 set(form_data.pop('skill_ids', None) or [])
2774 | set(model.get('info', {}).get('meta', {}).get('skillIds', []))
2775 | mentioned_skill_ids
2776 )
2777 available_skills = []
2778 terminal_skills = []
2779 view_skill_ids = []
2780 chat = None
2781 if is_saved_chat_id(metadata.get('chat_id')):
2782 chat = await Chats.get_chat_by_id(metadata['chat_id'])
2784 is_note_chat = bool(chat and (chat.meta or {}).get('internal') is True and (chat.meta or {}).get('type') == 'note')
2786 if is_note_chat:
2787 note_id = (chat.meta or {}).get('note_id')
2788 note = await Notes.get_note_by_id(note_id) if note_id else None
2789 if note and (
2790 user.role == 'admin'
2791 or note.user_id == user.id
2792 or await AccessGrants.has_access(
2793 user_id=user.id,
2794 resource_type='note',
2795 resource_id=note.id,
2796 permission='read',
2797 )
2798 ):
2799 note_files = [
2800 file
2801 for file in ((note.data or {}).get('files') or [])
2802 if isinstance(file, dict)
2803 and file.get('type') != 'image'
2804 and not (file.get('content_type') or '').startswith('image/')
2805 ]
2806 if note_files:
2807 files = [*(files or []), *note_files]
2809 use_builtin_tools = is_note_chat or (
2810 bool(metadata.get('session_id'))
2811 and metadata.get('params', {}).get('function_calling') != 'legacy'
2812 and (model.get('info', {}).get('meta', {}).get('capabilities') or {}).get('builtin_tools', True)
2813 )
2815 if skill_ids or use_builtin_tools:
2816 import aiohttp
2817 from open_webui.env import AIOHTTP_CLIENT_SESSION_TOOL_SERVER_SSL, AIOHTTP_CLIENT_TIMEOUT_TOOL_SERVER_DATA
2818 from open_webui.models.skills import Skills as SkillsModel
2819 from open_webui.utils.terminals import (
2820 format_terminal_skill_context,
2821 format_terminal_skill_manifest_entry,
2822 get_terminal_request_info,
2823 get_terminal_skill,
2824 )
2826 terminal_skill_prefix = 'terminal:'
2827 db_skill_ids = [sid for sid in skill_ids if not sid.startswith(terminal_skill_prefix)]
2828 terminal_skill_ids = [sid for sid in skill_ids if sid.startswith(terminal_skill_prefix)]
2830 if use_builtin_tools:
2831 accessible_skills = {s.id: s for s in await SkillsModel.get_skills(user_id=user.id)}
2832 db_skill_ids = sorted(accessible_skills)
2833 else:
2834 accessible_skills = {s.id: s for s in await SkillsModel.get_skills(user_id=user.id, ids=db_skill_ids)}
2836 for sid in db_skill_ids:
2837 skill = accessible_skills.get(sid)
2838 if skill and skill.is_active:
2839 available_skills.append(skill)
2841 skill_manifest = ''
2842 for skill in available_skills:
2843 if skill.id in mentioned_skill_ids or not use_builtin_tools:
2844 form_data['messages'] = add_or_update_system_message(
2845 f'<skill name="{skill.name}">\n{skill.content}\n</skill>',
2846 form_data['messages'],
2847 append=True,
2848 )
2849 else:
2850 view_skill_ids.append(skill.id)
2851 skill_manifest += (
2852 f'<skill>\n<id>{skill.id}</id>\n<name>{skill.name}</name>\n'
2853 f'<description>{skill.description or ""}</description>\n</skill>\n'
2854 )
2856 terminal_request = (
2857 await get_terminal_request_info(request, user, metadata, extra_params)
2858 if terminal_id or terminal_skill_ids
2859 else None
2860 )
2862 listed_terminal_skills = []
2863 if terminal_request:
2864 terminal_base_url, terminal_headers, terminal_cookies = terminal_request
2865 timeout = aiohttp.ClientTimeout(total=AIOHTTP_CLIENT_TIMEOUT_TOOL_SERVER_DATA)
2866 async with aiohttp.ClientSession(timeout=timeout, trust_env=True) as session:
2867 async with session.get(
2868 f'{terminal_base_url.rstrip("/")}/skills',
2869 headers=terminal_headers,
2870 cookies=terminal_cookies,
2871 ssl=AIOHTTP_CLIENT_SESSION_TOOL_SERVER_SSL,
2872 ) as response:
2873 if response.status == 200:
2874 listed = await response.json()
2875 listed_terminal_skills = listed if isinstance(listed, list) else []
2877 if terminal_id and use_builtin_tools:
2878 terminal_skills = listed_terminal_skills
2879 elif terminal_skill_ids:
2880 terminal_skill_map = {skill['id']: skill for skill in listed_terminal_skills}
2881 terminal_skills = [skill for sid in terminal_skill_ids if (skill := terminal_skill_map.get(sid))]
2883 for skill in terminal_skills:
2884 sid = skill['id']
2885 if sid in mentioned_skill_ids or not use_builtin_tools:
2886 skill_name = unquote(sid.removeprefix(terminal_skill_prefix))
2887 loaded = await get_terminal_skill(
2888 request, user.model_dump(), metadata, skill_name, extra_params
2889 )
2890 if loaded:
2891 form_data['messages'] = add_or_update_system_message(
2892 format_terminal_skill_context(loaded),
2893 form_data['messages'],
2894 append=True,
2895 )
2896 else:
2897 view_skill_ids.append(sid)
2898 skill_manifest += format_terminal_skill_manifest_entry(skill)
2900 if skill_manifest:
2901 form_data['messages'] = add_or_update_system_message(
2902 f'<available_skills>\n{skill_manifest}</available_skills>',
2903 form_data['messages'],
2904 append=True,
2905 )
2907 # Strip only resolved skill mentions; ordinary text such as Perl's <$fh> stays intact.
2908 resolved_skill_ids = {s.id for s in available_skills} | {s['id'] for s in terminal_skills}
2909 strip_skill_mentions(form_data.get('messages', []), resolved_skill_ids)
2911 prompt = get_last_user_message(form_data['messages'])
2913 # Guard against empty user message after skill mention stripping.
2914 # When a user selects a skill ($skill-name) without typing additional text,
2915 # the stripped result is an empty string which causes 400 errors on providers
2916 # that reject empty content blocks (e.g. AWS Bedrock ConverseStream).
2917 if not prompt or not prompt.strip():
2918 fallback = ', '.join([s.name for s in available_skills] + [s['name'] for s in terminal_skills])
2919 # Attachment-only messages keep their empty text, same as on models without skills.
2920 if fallback and not (metadata.get('user_message') or {}).get('files'):
2921 set_last_user_message_content(fallback, form_data['messages'])
2922 prompt = fallback
2923 # TODO: re-enable URL extraction from prompt
2924 # urls = []
2925 # if prompt and len(prompt or "") < 500 and (not files or len(files) == 0):
2926 # urls = extract_urls(prompt)
2928 if files:
2929 # files = [*files, *[{"type": "url", "url": url, "name": url} for url in urls]]
2930 # Remove duplicate files based on their content
2931 files = list({json.dumps(f, sort_keys=True): f for f in files}.values())
2933 metadata.update(
2934 {
2935 'model_id': form_data.get('model'),
2936 'tool_ids': tool_ids,
2937 'skill_ids': skill_ids,
2938 'terminal_id': terminal_id,
2939 'files': files,
2940 'features': features,
2941 }
2942 )
2943 form_data['metadata'] = metadata
2945 # When the caller provides an explicit `tools` key in the request body,
2946 # skip all server-side tool resolution and pass the caller's tools through
2947 # unchanged. Sending `tools: []` explicitly opts out of builtin injection.
2948 if payload_tools is None:
2949 # Server side tools
2950 tool_ids = metadata.get('tool_ids', None)
2951 # Client side tools
2952 direct_tool_servers = metadata.get('tool_servers', None)
2954 log.debug('tool_ids=%r', tool_ids)
2955 log.debug('direct_tool_servers=%r', direct_tool_servers)
2957 tools_dict = {}
2959 mcp_clients = {}
2960 mcp_tools_dict = {}
2962 if tool_ids:
2963 db_tool_ids = []
2964 for tool_id in tool_ids:
2965 if tool_id.startswith('server:mcp:'):
2966 try:
2967 server_id = tool_id[len('server:mcp:') :]
2969 result = await connect_mcp_server(
2970 request,
2971 server_id,
2972 user,
2973 metadata,
2974 extra_params,
2975 )
2976 if result is None:
2977 continue
2979 client, tool_specs = result
2980 mcp_clients[server_id] = client
2982 for tool_spec in tool_specs:
2984 async def make_tool_function(client, function_name):
2985 async def tool_function(**kwargs):
2986 return await client.call_tool(
2987 function_name,
2988 function_args=kwargs,
2989 )
2991 return tool_function
2993 tool_function = await make_tool_function(client, tool_spec['name'])
2995 mcp_tools_dict[f'{server_id}_{tool_spec["name"]}'] = {
2996 'spec': {
2997 **tool_spec,
2998 'name': f'{server_id}_{tool_spec["name"]}',
2999 },
3000 'callable': tool_function,
3001 'type': 'mcp',
3002 'client': client,
3003 'direct': False,
3004 }
3005 except Exception as e:
3006 log.debug(e)
3007 if event_emitter:
3008 await event_emitter(
3009 {
3010 'type': 'chat:message:error',
3011 'data': {'error': {'content': f"Failed to connect to MCP server '{server_id}'"}},
3012 }
3013 )
3014 continue
3015 elif ENABLE_PLUGINS:
3016 db_tool_ids.append(tool_id)
3018 if db_tool_ids:
3019 tools_dict = await get_tools(
3020 request,
3021 db_tool_ids,
3022 user,
3023 {
3024 **extra_params,
3025 '__model__': models[task_model_id],
3026 '__messages__': form_data['messages'],
3027 '__files__': metadata.get('files', []),
3028 },
3029 )
3031 if mcp_tools_dict:
3032 tools_dict = {**tools_dict, **mcp_tools_dict}
3034 # Resolve terminal tools if terminal_id is set (outside tool_ids check
3035 # so system terminals work even when no other tools are selected)
3036 terminal_capability = (model.get('info', {}).get('meta', {}).get('capabilities') or {}).get('terminal', True)
3037 terminal_connection_ids = {
3038 connection.get('id') for connection in await Config.get('terminal_server.connections', []) or []
3039 }
3040 if terminal_id and terminal_capability and terminal_id in terminal_connection_ids:
3041 try:
3042 terminal_result = await get_terminal_tools(
3043 request,
3044 terminal_id,
3045 user,
3046 extra_params,
3047 )
3048 if isinstance(terminal_result, tuple):
3049 terminal_tools, system_prompt = terminal_result
3050 else:
3051 terminal_tools = terminal_result
3052 system_prompt = None
3053 if terminal_tools:
3054 tools_dict = {**tools_dict, **terminal_tools}
3055 if system_prompt:
3056 form_data['messages'] = add_or_update_system_message(
3057 system_prompt,
3058 form_data['messages'],
3059 append=True,
3060 )
3061 except Exception as e:
3062 log.exception(e)
3063 raise HTTPException(status_code=503, detail=f'Terminal unavailable: {e}') from e
3065 if direct_tool_servers:
3066 for tool_server in direct_tool_servers:
3067 if tool_server.get('is_terminal') is True and not terminal_capability:
3068 continue
3069 system_prompt = tool_server.pop('system_prompt', None)
3070 if system_prompt:
3071 form_data['messages'] = add_or_update_system_message(
3072 system_prompt,
3073 form_data['messages'],
3074 append=True,
3075 )
3077 tool_specs = tool_server.pop('specs', [])
3079 for tool in tool_specs:
3080 tools_dict[tool['name']] = {
3081 'spec': tool,
3082 'direct': True,
3083 'server': tool_server,
3084 }
3086 if terminal_id and terminal_capability:
3087 from open_webui.utils.terminals import add_terminal_agents_md, get_terminal_agents_md
3089 agents_md = await get_terminal_agents_md(request, user, metadata, extra_params)
3090 if agents_md:
3091 form_data['messages'] = add_terminal_agents_md(form_data['messages'], agents_md)
3093 if mcp_clients:
3094 metadata['mcp_clients'] = mcp_clients
3096 # Inject builtin tools for native function calling based on enabled features and model capability.
3097 # Only inject when the request originates from the UI (identified by session_id).
3098 # API callers don't expect hidden tools; they can explicitly request tools via tool_ids.
3099 if use_builtin_tools:
3100 # Add file context to user messages
3101 chat_id = metadata.get('chat_id')
3102 form_data['messages'] = await add_file_context(form_data.get('messages', []), chat_id, user)
3104 if (model.get('info', {}).get('meta', {}).get('builtinTools') or {}).get('knowledge', True):
3105 from html import escape
3107 knowledge_tags = []
3108 for item in get_attached_knowledge(model, metadata):
3109 if not item.get('id') or not item.get('type'):
3110 continue
3111 attrs = f'type="{escape(str(item["type"]), quote=True)}" id="{escape(str(item["id"]), quote=True)}"'
3112 if item.get('name'):
3113 attrs += f' name="{escape(str(item["name"]), quote=True)}"'
3114 if item.get('source'):
3115 attrs += f' source="{escape(str(item["source"]), quote=True)}"'
3116 knowledge_tags.append(f'<knowledge {attrs}/>')
3118 if knowledge_tags:
3119 form_data['messages'] = add_or_update_system_message(
3120 '<attached_knowledge>\n' + '\n'.join(knowledge_tags) + '\n</attached_knowledge>',
3121 form_data['messages'],
3122 append=True,
3123 )
3125 builtin_tools = await get_builtin_tools(
3126 request,
3127 {
3128 **extra_params,
3129 '__event_emitter__': event_emitter,
3130 '__skill_ids__': view_skill_ids,
3131 },
3132 features,
3133 model,
3134 is_note_chat=is_note_chat,
3135 )
3136 for name, tool_dict in builtin_tools.items():
3137 if name not in tools_dict:
3138 tools_dict[name] = tool_dict
3140 # Only advertise user-shell tools when the originating browser has a connected shell.
3141 shell_tools = {
3142 name: tool
3143 for name, tool in tools_dict.items()
3144 if name in {'read_user_terminal', 'send_user_terminal_input'}
3145 and (tool.get('type') == 'terminal' or tool.get('server', {}).get('is_terminal') is True)
3146 }
3147 selected = {
3148 name
3149 for name, tool in shell_tools.items()
3150 if terminal_id
3151 and (
3152 tool.get('tool_id') == f'terminal:{terminal_id}'
3153 or (tool.get('direct') and tool.get('server', {}).get('url') == terminal_id)
3154 )
3155 }
3156 connected = False
3157 if (
3158 selected
3159 and event_caller
3160 and metadata.get('session_id')
3161 and metadata.get('chat_id')
3162 and not metadata.get('automation_id')
3163 and not metadata.get('internal')
3164 ):
3165 try:
3166 state = await asyncio.wait_for(
3167 event_caller(
3168 {
3169 'type': 'request:terminal:state',
3170 'data': {'terminal_id': terminal_id, 'session_id': metadata['session_id']},
3171 }
3172 ),
3173 timeout=2,
3174 )
3175 connected = isinstance(state, dict) and state.get('connected') is True
3176 except Exception:
3177 # Old/disconnected browsers cannot confirm availability; other tools still work.
3178 pass
3180 for name in shell_tools:
3181 if not connected or name not in selected:
3182 tools_dict.pop(name)
3184 if tools_dict:
3185 # Always store resolved tools in metadata so downstream consumers
3186 # (e.g. pipe functions) can access all tools including MCP and builtins.
3187 metadata['tools'] = tools_dict
3189 if metadata.get('params', {}).get('function_calling') != 'legacy':
3190 # If the function calling is native, then call the tools function calling handler
3191 form_data['tools'] = [
3192 {'type': 'function', 'function': tool.get('spec', {})} for tool in tools_dict.values()
3193 ]
3194 if inlet_filter_tools:
3195 form_data['tools'].extend(inlet_filter_tools)
3196 else:
3197 # If the function calling is not native, then call the tools function calling handler
3198 try:
3199 form_data, flags = await chat_completion_tools_handler(
3200 request, form_data, extra_params, user, models, tools_dict
3201 )
3202 sources.extend(flags.get('sources', []))
3203 except Exception as e:
3204 log.exception(e)
3206 # Check if file context extraction is enabled for this model (default True)
3207 file_context_enabled = (model.get('info', {}).get('meta', {}).get('capabilities') or {}).get('file_context', True)
3209 if file_context_enabled:
3210 try:
3211 form_data, flags = await chat_completion_files_handler(request, form_data, extra_params, user)
3212 sources.extend(flags.get('sources', []))
3213 except Exception as e:
3214 log.exception(e)
3216 # Save the pre-RAG message state so the native tool call loop can
3217 # restore to the true original (before file-source injection) rather
3218 # than a snapshot that already has the RAG template baked in.
3219 system_message = get_system_message(form_data['messages'])
3220 system_content = get_content_from_message(system_message) if system_message else ''
3221 resolved_model_system_prompt = await resolve_system_prompt(
3222 model_system_prompt,
3223 metadata,
3224 user,
3225 )
3226 if resolved_model_system_prompt:
3227 system_content = (
3228 f'{resolved_model_system_prompt}\n{system_content}' if system_content else resolved_model_system_prompt
3229 )
3230 metadata['system_prompt'] = system_content or None
3231 metadata['user_prompt'] = get_last_user_message(form_data['messages'])
3232 metadata['sources'] = sources[:] if sources else []
3234 # If context is not empty, insert it into the messages
3235 if sources and prompt:
3236 form_data['messages'] = await apply_source_context_to_messages(request, form_data['messages'], sources, prompt)
3238 # If there are citations, add them to the data_items
3239 sources = [
3240 source
3241 for source in sources
3242 if source.get('source', {}).get('name', '') or source.get('source', {}).get('id', '')
3243 ]
3245 if len(sources) > 0:
3246 events.append({'sources': sources})
3248 if model_knowledge:
3249 await event_emitter(
3250 {
3251 'type': 'status',
3252 'data': {
3253 'action': 'knowledge_search',
3254 'query': user_message,
3255 'done': True,
3256 'hidden': True,
3257 },
3258 }
3259 )
3261 if ENABLE_PLUGINS:
3262 try:
3263 form_data, _ = await process_filter_functions(
3264 request=request,
3265 filter_context=filter_context,
3266 filter_functions=filter_functions,
3267 filter_type='request',
3268 form_data=form_data,
3269 extra_params=extra_params,
3270 )
3271 except Exception as e:
3272 raise Exception(f'{e}')
3274 form_data = normalize_messages_for_model(form_data)
3276 return form_data, metadata, events
3279async def get_event_emitter_and_caller(metadata):
3280 event_emitter = None
3281 event_caller = None
3283 # event_emitter only needs user_id + chat_id + message_id.
3284 # It broadcasts to user:{user_id} room AND persists to DB,
3285 # so it works for backend-initiated calls (automations, API).
3286 if metadata.get('chat_id') and metadata.get('message_id'):
3287 event_emitter = await get_event_emitter(metadata)
3289 # event_caller needs session_id — it calls back to a specific
3290 # websocket session (used by direct tools, pyodide code interpreter).
3291 if metadata.get('session_id') and metadata.get('chat_id') and metadata.get('message_id'):
3292 event_caller = await get_event_call(metadata)
3294 return event_emitter, event_caller
3297async def build_chat_response_context(request, form_data, user, model, metadata, tasks, events):
3298 event_emitter, event_caller = await get_event_emitter_and_caller(metadata)
3299 assistant_message = None
3300 if metadata.get('assistant_message_id'):
3301 # Preserve the original output before request conversion strips it from temporary chats.
3302 if is_saved_chat_id(metadata.get('chat_id')):
3303 assistant_message = await Chats.get_message_by_id_and_message_id(
3304 metadata['chat_id'], metadata['assistant_message_id']
3305 )
3306 elif form_data.get('messages') and form_data['messages'][-1].get('role') == 'assistant':
3307 assistant_message = form_data['messages'][-1]
3308 return {
3309 'request': request,
3310 'form_data': form_data,
3311 'user': user,
3312 'model': model,
3313 'metadata': metadata,
3314 'tasks': tasks,
3315 'events': events,
3316 'event_emitter': event_emitter,
3317 'event_caller': event_caller,
3318 'assistant_message': copy.deepcopy(assistant_message),
3319 }
3322async def execute_tool_call_for_output(request, form_data, user, metadata, event_caller, event_emitter, tool_call):
3323 tools = metadata.get('tools', {})
3324 name = tool_call.get('function', {}).get('name', '')
3325 tool_args = tool_call.get('function', {}).get('arguments', '{}')
3326 params = {}
3327 if tool_args and tool_args.strip():
3328 try:
3329 params = JSONCodec.loads(tool_args)
3330 except Exception:
3331 try:
3332 params = ast.literal_eval(tool_args)
3333 except Exception as e:
3334 log.debug(e)
3335 return {
3336 'tool_call_id': tool_call.get('id', ''),
3337 'content': (
3338 'Error: Tool call arguments could not be parsed. '
3339 'The model generated malformed or incomplete JSON.'
3340 ),
3341 }
3342 if not isinstance(params, dict):
3343 return {
3344 'tool_call_id': tool_call.get('id', ''),
3345 'content': 'Error: Tool call arguments must be a JSON object.',
3346 }
3347 tool_call.setdefault('function', {})['arguments'] = JSONCodec.dumps(params)
3349 tool = tools.get(name)
3350 if not tool:
3351 return {'tool_call_id': tool_call.get('id', ''), 'content': f'Error: Tool "{name}" not found.'}
3353 spec = tool.get('spec', {})
3354 tool_type = tool.get('type', '')
3355 direct_tool = tool.get('direct', False)
3356 allowed_params = spec.get('parameters', {}).get('properties', {}).keys()
3357 params = {key: value for key, value in params.items() if key in allowed_params}
3359 try:
3360 if direct_tool:
3361 if not event_caller:
3362 result = 'Error: Browser session is not connected for this direct tool.'
3363 else:
3364 result = await event_caller(
3365 {
3366 'type': 'execute:tool',
3367 'data': {
3368 'id': str(uuid4()),
3369 'name': name,
3370 'params': params,
3371 'server': tool.get('server', {}),
3372 'session_id': metadata.get('session_id'),
3373 },
3374 }
3375 )
3376 else:
3377 function = await get_updated_tool_function(
3378 function=tool['callable'],
3379 extra_params={
3380 '__messages__': form_data.get('messages', []),
3381 '__files__': metadata.get('files', []),
3382 },
3383 )
3384 result = await function(**params)
3385 except Exception as e:
3386 result = {'error': str(e)}
3388 terminal_file_result = build_terminal_file_tool_result(name, params, result, tool, metadata)
3389 if terminal_file_result:
3390 result = terminal_file_result
3392 result, files, embeds = await process_tool_result(
3393 request,
3394 name,
3395 result,
3396 tool_type,
3397 direct_tool,
3398 metadata,
3399 user,
3400 )
3402 await terminal_event_handler(name, params, result, event_emitter)
3404 return {
3405 'tool_call_id': tool_call.get('id', ''),
3406 'content': tool_result_content(result),
3407 **({'files': files} if files else {}),
3408 **({'embeds': embeds} if embeds else {}),
3409 }
3412async def drain_approved_tool_calls(request, form_data, user, model, metadata) -> bool:
3413 chat_id = metadata.get('chat_id')
3414 assistant_message_id = metadata.get('assistant_message_id')
3415 # Only a resume/continue payload re-enters an existing message; other paths mint a fresh id with nothing to drain.
3416 if not is_saved_chat_id(chat_id) or not assistant_message_id:
3417 return False
3419 message_id = metadata.get('message_id') or assistant_message_id
3420 message = await Chats.get_message_by_id_and_message_id(chat_id, message_id)
3421 output = message.get('output') if message else None
3422 if not isinstance(output, list):
3423 return False
3425 result_call_ids = {
3426 item.get('call_id') for item in output if item.get('type') == 'function_call_output' and item.get('call_id')
3427 }
3428 approved_calls = [
3429 item
3430 for item in output
3431 if item.get('type') == 'function_call'
3432 and item.get('call_id')
3433 and item.get('status') == 'queued'
3434 and item.get('approved') is True
3435 and item.get('call_id') not in result_call_ids
3436 ]
3437 if not approved_calls:
3438 if metadata.get('params', {}).get('tool_approval_mode', 'full') == 'ask' and any(
3439 item.get('type') == 'function_call'
3440 and item.get('name') != 'ask_user'
3441 and (item.get('call_id') or item.get('id'))
3442 and item.get('status') == 'queued'
3443 and item.get('approved') is not True
3444 and (item.get('call_id') or item.get('id')) not in result_call_ids
3445 for item in output
3446 ):
3447 event_emitter, _ = await get_event_emitter_and_caller(metadata)
3448 await pause_for_tool_approval(chat_id, message_id, output, form_data, metadata)
3449 if event_emitter:
3450 await event_emitter({'type': 'chat:completion', 'data': {'done': False, 'output': output}})
3451 return True
3452 return False
3454 event_emitter, event_caller = await get_event_emitter_and_caller(metadata)
3455 changed = False
3456 for item in approved_calls:
3457 if item.get('name') == 'ask_user':
3458 item['status'] = 'pending'
3459 item.pop('approved', None)
3460 changed = True
3461 continue
3463 tool_call = {
3464 'id': item.get('call_id', ''),
3465 'type': 'function',
3466 'function': {
3467 'name': item.get('name', ''),
3468 'arguments': item.get('arguments', '{}'),
3469 },
3470 }
3471 result = await execute_tool_call_for_output(
3472 request,
3473 form_data,
3474 user,
3475 metadata,
3476 event_caller,
3477 event_emitter,
3478 tool_call,
3479 )
3480 item['arguments'] = tool_call.get('function', {}).get('arguments', '{}')
3481 output_parts = [{'type': 'input_text', 'text': result.get('content', '')}]
3482 item['status'] = 'failed' if _is_tool_result_error(result.get('content', '')) else 'completed'
3483 display_files = []
3484 for file_item in result.get('files', []):
3485 if file_item.get('type') == 'image' and file_item.get('url', '').startswith('data:'):
3486 image_url = await store_tool_result_image(request, file_item['url'], metadata, user)
3487 output_parts.append({'type': 'input_image', 'image_url': image_url})
3488 else:
3489 display_files.append(file_item)
3491 output.append(
3492 {
3493 'type': 'function_call_output',
3494 'id': output_id('fco'),
3495 'call_id': result.get('tool_call_id', ''),
3496 'output': output_parts,
3497 'status': item['status'],
3498 **({'files': display_files} if display_files else {}),
3499 **({'embeds': result.get('embeds')} if result.get('embeds') else {}),
3500 }
3501 )
3502 changed = True
3504 if changed:
3505 result_call_ids = {
3506 item.get('call_id') for item in output if item.get('type') == 'function_call_output' and item.get('call_id')
3507 }
3508 if metadata.get('params', {}).get('tool_approval_mode', 'full') == 'ask' and any(
3509 item.get('type') == 'function_call'
3510 and item.get('name') != 'ask_user'
3511 and (item.get('call_id') or item.get('id'))
3512 and item.get('status') == 'queued'
3513 and item.get('approved') is not True
3514 and (item.get('call_id') or item.get('id')) not in result_call_ids
3515 for item in output
3516 ):
3517 await pause_for_tool_approval(chat_id, message_id, output, form_data, metadata)
3518 result_call_ids = {
3519 item.get('call_id')
3520 for item in output
3521 if item.get('type') == 'function_call_output' and item.get('call_id')
3522 }
3523 paused = any(
3524 item.get('type') == 'function_call'
3525 and item.get('call_id')
3526 and item.get('status') in {'pending', 'queued', 'requires_approval'}
3527 and item.get('call_id') not in result_call_ids
3528 for item in output
3529 )
3530 if not paused:
3531 output.append(
3532 {
3533 'type': 'message',
3534 'id': output_id('msg'),
3535 'status': 'in_progress',
3536 'role': 'assistant',
3537 'content': [{'type': 'output_text', 'text': ''}],
3538 }
3539 )
3541 await Chats.upsert_message_to_chat_by_id_and_message_id(
3542 chat_id,
3543 message_id,
3544 {'done': False, 'output': output},
3545 touch=False,
3546 )
3547 if event_emitter:
3548 await event_emitter(
3549 {
3550 'type': 'chat:completion',
3551 'data': {
3552 'done': False,
3553 'output': output,
3554 },
3555 }
3556 )
3558 db_messages = await load_messages_from_db(chat_id, metadata.get('user_message_id'))
3559 if db_messages:
3560 assistant_message = await Chats.get_message_by_id_and_message_id(chat_id, message_id)
3561 if assistant_message:
3562 db_messages.append({k: v for k, v in assistant_message.items() if k in MESSAGE_REPLAY_KEYS})
3563 for message in db_messages:
3564 output = message.get('output')
3565 # reasoning_details can be model/provider-bound, so only replay them
3566 # for output produced by the same model.
3567 if (
3568 message.get('role') == 'assistant'
3569 and message.get('model') != model['id']
3570 and isinstance(output, list)
3571 ):
3572 message['output'] = strip_reasoning_details(output)
3574 form_data['messages'] = process_messages_with_output(
3575 db_messages,
3576 reasoning_format=get_reasoning_format(model),
3577 )
3578 form_data['messages'] = sanitize_tool_pairs(form_data['messages'])
3580 if not paused and ENABLE_PLUGINS:
3581 filter_functions = await get_filter_functions(request, model, metadata.get('filter_ids', []))
3582 if filter_functions:
3583 filtered_form_data, _ = await process_filter_functions(
3584 request=request,
3585 filter_context=get_filter_context(request),
3586 filter_functions=filter_functions,
3587 filter_type='request',
3588 form_data=form_data,
3589 extra_params={
3590 '__event_emitter__': event_emitter,
3591 '__event_call__': event_caller,
3592 '__user__': user.model_dump() if isinstance(user, UserModel) else {},
3593 '__metadata__': metadata,
3594 '__oauth_token__': await get_system_oauth_token(request, user),
3595 '__request__': request,
3596 '__model__': model,
3597 '__chat_id__': metadata.get('chat_id'),
3598 '__message_id__': metadata.get('message_id'),
3599 },
3600 )
3601 if filtered_form_data is not form_data:
3602 form_data.clear()
3603 form_data.update(filtered_form_data)
3605 if not paused:
3606 normalize_messages_for_model(form_data)
3608 return paused
3610 return False
3613async def pause_for_tool_approval(chat_id: str, message_id: str, output: list[dict], form_data: dict, metadata: dict):
3614 result_call_ids = {
3615 item.get('call_id') for item in output if item.get('type') == 'function_call_output' and item.get('call_id')
3616 }
3617 has_pending_approval = False
3618 for item in output:
3619 if item.get('type') == 'function_call' and not item.get('call_id') and item.get('id'):
3620 item['call_id'] = item['id']
3622 if (
3623 item.get('type') == 'function_call'
3624 and item.get('call_id')
3625 and item.get('call_id') not in result_call_ids
3626 and item.get('status') != 'rejected'
3627 ):
3628 if not has_pending_approval:
3629 item['status'] = 'pending'
3630 has_pending_approval = True
3631 elif item.get('status') == 'in_progress':
3632 item['status'] = 'queued'
3634 await Chats.upsert_message_to_chat_by_id_and_message_id(
3635 chat_id,
3636 message_id,
3637 {
3638 'done': False,
3639 'output': output,
3640 'meta': {
3641 **(metadata.get('tool_approval') or {}),
3642 'session_id': metadata.get('session_id'),
3643 'tool_ids': metadata.get('tool_ids') or [],
3644 'skill_ids': metadata.get('skill_ids') or [],
3645 'terminal_id': metadata.get('terminal_id'),
3646 'tool_servers': metadata.get('tool_servers'),
3647 'filter_ids': metadata.get('filter_ids') or [],
3648 'features': metadata.get('features') or {},
3649 'variables': metadata.get('variables') or {},
3650 'files': metadata.get('files') or [],
3651 'params': metadata.get('params') or {},
3652 },
3653 },
3654 touch=False,
3655 )
3658def get_response_data(response):
3659 if isinstance(response, list) and len(response) == 1:
3660 # If the response is a single-item list, unwrap it #17213
3661 response = response[0]
3663 if isinstance(response, JSONResponse):
3664 if isinstance(response.body, bytes):
3665 try:
3666 response_data = JSONCodec.loads(response.body.decode('utf-8', 'replace'))
3667 except JSONCodec.JSONDecodeError:
3668 response_data = {'error': {'detail': 'Invalid JSON response'}}
3669 else:
3670 response_data = response
3671 elif isinstance(response, dict):
3672 response_data = response
3673 else:
3674 response_data = None
3676 return response, response_data
3679def merge_events_into_response(response_data, events):
3680 if events and isinstance(events, list):
3681 extra_response = {}
3682 for event in events:
3683 if isinstance(event, dict):
3684 extra_response.update(event)
3685 else:
3686 extra_response[event] = True
3688 return {
3689 **extra_response,
3690 **response_data,
3691 }
3692 return response_data
3695def build_response_object(response, response_data):
3696 if isinstance(response, dict):
3697 return response_data
3698 if isinstance(response, JSONResponse):
3699 return JSONResponse(
3700 content=response_data,
3701 headers=response.headers,
3702 status_code=response.status_code,
3703 )
3704 return response
3707def update_assistant_message_from_stream(assistant_message, raw):
3708 line = raw.decode('utf-8', 'replace') if isinstance(raw, bytes) else raw
3709 if not isinstance(line, str):
3710 return
3712 def append_output_text(item, text):
3713 parts = item.setdefault('content', [])
3714 if parts and parts[-1].get('type') == 'output_text':
3715 append_to_text_field(parts[-1], 'text', text)
3716 else:
3717 parts.append({'type': 'output_text', 'text': text})
3719 for raw_part in line.splitlines():
3720 part = raw_part.removeprefix('data:').strip()
3721 if not part or part == '[DONE]':
3722 continue
3724 try:
3725 data = JSONCodec.loads(part)
3726 except Exception:
3727 continue
3729 if not isinstance(data, dict):
3730 continue
3732 if data.get('type', '').startswith('response.'):
3733 output, meta = handle_responses_streaming_event(data, assistant_message.get('output', []))
3734 if output:
3735 assistant_message['output'] = output
3736 if meta and meta.get('usage'):
3737 assistant_message['usage'] = merge_usage(assistant_message.get('usage'), meta['usage'])
3738 continue
3740 raw_usage = data.get('usage', {}) or {}
3741 raw_usage.update(data.get('timings', {}))
3742 if raw_usage:
3743 assistant_message['usage'] = merge_usage(assistant_message.get('usage'), raw_usage)
3745 for choice in data.get('choices', []):
3746 delta = choice.get('delta', {}) or {}
3747 content = delta.get('content')
3748 reasoning_content = delta.get('reasoning_content') or delta.get('reasoning') or delta.get('thinking')
3750 if reasoning_content:
3751 output = assistant_message.setdefault('output', [])
3752 if not output or output[-1].get('type') != 'reasoning':
3753 output.append(
3754 {
3755 'type': 'reasoning',
3756 'id': output_id('r'),
3757 'status': 'in_progress',
3758 'start_tag': '<think>',
3759 'end_tag': '</think>',
3760 'attributes': {'type': 'reasoning_content'},
3761 'content': [],
3762 'summary': None,
3763 'started_at': time.time(),
3764 }
3765 )
3767 append_output_text(output[-1], reasoning_content)
3769 if content:
3770 output = assistant_message.get('output')
3771 if output:
3772 if output[-1].get('type') == 'reasoning':
3773 output[-1]['status'] = 'completed'
3774 output[-1]['ended_at'] = time.time()
3775 output[-1]['duration'] = int(output[-1]['ended_at'] - output[-1]['started_at'])
3777 if not output or output[-1].get('type') != 'message':
3778 output.append(
3779 {
3780 'type': 'message',
3781 'id': output_id('msg'),
3782 'status': 'in_progress',
3783 'role': 'assistant',
3784 'content': [],
3785 }
3786 )
3788 append_output_text(output[-1], content)
3790 if 'content' in assistant_message:
3791 append_to_text_field(assistant_message, 'content', content)
3792 else:
3793 assistant_message['content'] = '' + content
3796async def get_system_oauth_token(request, user):
3797 """Get the system OAuth token for a user.
3799 Primary path: use the oauth_session_id cookie (browser requests).
3800 Fallback: look up the user's most recent OAuth session from the DB
3801 (covers automations, API calls, and other cookie-less contexts).
3802 """
3803 oauth_token = None
3804 try:
3805 oauth_session_id = request.cookies.get('oauth_session_id', None)
3806 if oauth_session_id:
3807 oauth_token = await request.app.state.oauth_manager.get_oauth_token(
3808 user.id,
3809 oauth_session_id,
3810 )
3812 # Fallback: no cookie (automation, API key, etc.) — use most recent session
3813 if oauth_token is None:
3814 from open_webui.models.oauth_sessions import OAuthSessions
3816 sessions = await OAuthSessions.get_sessions_by_user_id(user.id)
3817 # Filter out MCP-provider sessions — their token refresh is handled
3818 # separately by oauth_client_manager. Passing them to the SSO
3819 # oauth_manager causes a failed refresh and session deletion (#24618).
3820 sessions = [s for s in sessions if not (s.provider or '').startswith('mcp:')]
3821 if sessions:
3822 best = max(sessions, key=lambda s: s.updated_at)
3823 oauth_token = await request.app.state.oauth_manager.get_oauth_token(
3824 user.id,
3825 best.id,
3826 )
3827 except Exception as e:
3828 log.error(f'Error getting OAuth token: {e}')
3829 return oauth_token
3832async def background_tasks_handler(ctx):
3833 request = ctx['request']
3834 form_data = ctx['form_data']
3835 user = ctx['user']
3836 metadata = ctx['metadata']
3837 tasks = ctx['tasks']
3838 event_emitter = ctx['event_emitter']
3840 message = None
3841 messages = []
3843 if is_saved_chat_id(metadata.get('chat_id')):
3844 messages_map = await Chats.get_messages_map_by_chat_id(metadata['chat_id'])
3845 if not messages_map:
3846 # Chat was deleted while the response was streaming — skip background tasks
3847 return
3848 message = messages_map.get(metadata['message_id'])
3850 message_list = get_message_list(messages_map, metadata['message_id'])
3852 # Remove details tags and files from the messages.
3853 # as get_message_list creates a new list, it does not affect
3854 # the original messages outside of this handler
3856 messages = []
3857 for message in message_list:
3858 content = message.get('content', '')
3859 if isinstance(content, list):
3860 for item in content:
3861 if item.get('type') == 'text':
3862 content = item['text']
3863 break
3865 if isinstance(content, str):
3866 content = re.sub(
3867 r'<details\b[^>]*>.*?<\/details>|!\[.*?\]\(.*?\)',
3868 '',
3869 content,
3870 flags=re.S | re.I,
3871 ).strip()
3873 messages.append(
3874 {
3875 **message,
3876 'role': message.get('role', 'assistant'), # Safe fallback for missing role
3877 'content': content,
3878 }
3879 )
3880 else:
3881 # Local temp chat, get the model and message from the form_data
3882 message = get_last_user_message_item(form_data.get('messages', []))
3883 messages = form_data.get('messages', [])
3884 if message:
3885 message['model'] = form_data.get('model')
3887 if message and 'model' in message:
3888 if tasks and messages:
3889 if TASKS.FOLLOW_UP_GENERATION in tasks and tasks[TASKS.FOLLOW_UP_GENERATION]:
3890 res = await generate_follow_ups(
3891 request,
3892 {
3893 'model': message['model'],
3894 'messages': messages,
3895 'message_id': metadata['message_id'],
3896 'chat_id': metadata['chat_id'],
3897 },
3898 user,
3899 )
3901 if res and isinstance(res, dict):
3902 if len(res.get('choices', [])) == 1:
3903 response_message = res.get('choices', [])[0].get('message', {})
3905 follow_ups_string = response_message.get('content') or response_message.get(
3906 'reasoning_content', ''
3907 )
3908 else:
3909 follow_ups_string = ''
3911 follow_ups_string = follow_ups_string[
3912 follow_ups_string.find('{') : follow_ups_string.rfind('}') + 1
3913 ]
3915 try:
3916 follow_ups = JSONCodec.loads(follow_ups_string).get('follow_ups', [])
3917 await event_emitter(
3918 {
3919 'type': 'chat:message:follow_ups',
3920 'data': {
3921 'follow_ups': follow_ups,
3922 },
3923 }
3924 )
3926 if is_saved_chat_id(metadata.get('chat_id')):
3927 await Chats.upsert_message_to_chat_by_id_and_message_id(
3928 metadata['chat_id'],
3929 metadata['message_id'],
3930 {
3931 'followUps': follow_ups,
3932 },
3933 touch=False,
3934 )
3936 except Exception as e:
3937 pass
3939 if is_saved_chat_id(metadata.get('chat_id')): # Only update titles and tags for saved chats
3940 if TASKS.TITLE_GENERATION in tasks:
3941 user_message = get_last_user_message(messages)
3942 if user_message and len(user_message) > 100:
3943 user_message = user_message[:100] + '...'
3945 title = None
3946 if tasks[TASKS.TITLE_GENERATION]:
3947 res = await generate_title(
3948 request,
3949 {
3950 'model': message['model'],
3951 'messages': messages,
3952 'chat_id': metadata['chat_id'],
3953 },
3954 user,
3955 )
3957 if res and isinstance(res, dict):
3958 if len(res.get('choices', [])) == 1:
3959 response_message = res.get('choices', [])[0].get('message', {})
3961 title_string = (
3962 response_message.get('content')
3963 or response_message.get(
3964 'reasoning_content',
3965 )
3966 or message.get('content', user_message)
3967 )
3968 else:
3969 title_string = ''
3971 title_string = title_string[title_string.find('{') : title_string.rfind('}') + 1]
3973 try:
3974 title = JSONCodec.loads(title_string).get('title', user_message)
3975 except Exception as e:
3976 title = ''
3978 if not title:
3979 title = messages[0].get('content', user_message)
3981 await Chats.update_chat_title_by_id(metadata['chat_id'], title)
3983 await event_emitter(
3984 {
3985 'type': 'chat:title',
3986 'data': title,
3987 }
3988 )
3990 if title == None and len(messages) == 2 and (not messages_map or len(messages_map) <= 2):
3991 title = messages[0].get('content', user_message)
3993 await Chats.update_chat_title_by_id(metadata['chat_id'], title)
3995 await event_emitter(
3996 {
3997 'type': 'chat:title',
3998 'data': message.get('content', user_message),
3999 }
4000 )
4002 if TASKS.TAGS_GENERATION in tasks and tasks[TASKS.TAGS_GENERATION]:
4003 res = await generate_chat_tags(
4004 request,
4005 {
4006 'model': message['model'],
4007 'messages': messages,
4008 'chat_id': metadata['chat_id'],
4009 },
4010 user,
4011 )
4013 if res and isinstance(res, dict):
4014 if len(res.get('choices', [])) == 1:
4015 response_message = res.get('choices', [])[0].get('message', {})
4017 tags_string = response_message.get('content') or response_message.get(
4018 'reasoning_content', ''
4019 )
4020 else:
4021 tags_string = ''
4023 tags_string = tags_string[tags_string.find('{') : tags_string.rfind('}') + 1]
4025 try:
4026 tags = JSONCodec.loads(tags_string).get('tags', [])
4027 await Chats.update_chat_tags_by_id(metadata['chat_id'], tags, user)
4029 await event_emitter(
4030 {
4031 'type': 'chat:tags',
4032 'data': tags,
4033 }
4034 )
4035 except Exception as e:
4036 pass
4038 if messages:
4039 await review_memory_after_turn(
4040 request=request,
4041 user=user,
4042 model=ctx['model'],
4043 metadata=metadata,
4044 form_data=form_data,
4045 assistant_message=ctx.get('assistant_message') or {},
4046 messages=messages,
4047 )
4050async def outlet_filter_handler(ctx):
4051 """Run outlet filters inline after chat completion.
4053 Replaces the separate POST /api/chat/completed round-trip.
4054 Persists outlet-modified content to DB and emits a chat:outlet event
4055 so the frontend can sync its in-memory state. Returns immediately when
4056 the model has no filters.
4058 For temp/API chats, messages are built from form_data plus ctx['assistant_message'].
4059 """
4060 request = ctx['request']
4061 user = ctx['user']
4062 model = ctx['model']
4063 metadata = ctx['metadata']
4064 event_emitter = ctx.get('event_emitter')
4065 event_caller = ctx.get('event_caller')
4067 chat_id = metadata.get('chat_id', '')
4068 message_id = metadata.get('message_id')
4070 if not chat_id and not ctx.get('assistant_message'):
4071 return
4073 if not message_id:
4074 message_id = output_id('msg')
4076 is_unsaved_chat = not is_saved_chat_id(chat_id)
4077 try:
4078 filter_functions = (
4079 await get_filter_functions(request, model, metadata.get('filter_ids', [])) if ENABLE_PLUGINS else []
4080 )
4081 model_id = model.get('id') if isinstance(model, dict) else model
4082 models = request.app.state.MODELS
4083 has_pipeline_outlet_filters = bool(
4084 (isinstance(model, dict) and 'pipeline' in model) or get_sorted_filters(model_id, models)
4085 )
4086 if not filter_functions and not has_pipeline_outlet_filters:
4087 return
4089 messages_map = None
4091 if is_unsaved_chat:
4092 form_messages = ctx.get('form_data', {}).get('messages', [])
4093 assistant_message = ctx.get('assistant_message', {})
4095 message_list = [
4096 {
4097 'role': m.get('role'),
4098 'content': m.get('content') or get_output_text(m.get('output')),
4099 }
4100 for m in form_messages
4101 ]
4103 if assistant_message:
4104 message_list.append(
4105 {
4106 'id': message_id,
4107 'role': 'assistant',
4108 **assistant_message,
4109 }
4110 )
4112 if not message_list:
4113 return
4114 else:
4115 messages_map = await Chats.get_messages_map_by_chat_id(chat_id)
4116 if not messages_map:
4117 return
4119 message_list = get_message_list(messages_map, message_id)
4120 if not message_list:
4121 return
4123 outlet_data = {
4124 'model': model_id,
4125 'messages': [
4126 {
4127 'id': m.get('id'),
4128 'role': m.get('role'),
4129 'content': m.get('content') or get_output_text(m.get('output')),
4130 'info': m.get('info'),
4131 'timestamp': m.get('timestamp'),
4132 # Deepcopy so in-place filter mutations do not alias messages_map's baseline
4133 **({'output': copy.deepcopy(m['output'])} if m.get('output') else {}),
4134 **({'usage': m['usage']} if m.get('usage') else {}),
4135 **({'sources': m['sources']} if m.get('sources') else {}),
4136 }
4137 for m in message_list
4138 ],
4139 'filter_ids': metadata.get('filter_ids', []),
4140 'chat_id': chat_id,
4141 'session_id': metadata.get('session_id'),
4142 'id': message_id,
4143 }
4145 # Pipeline outlet filters
4146 try:
4147 outlet_data = await process_pipeline_outlet_filter(request, outlet_data, user, models)
4148 except Exception as e:
4149 log.debug('Pipeline outlet filter error: %s', e)
4151 # Function outlet filters
4152 extra_params = {
4153 '__event_emitter__': event_emitter,
4154 '__event_call__': event_caller,
4155 '__user__': user.model_dump() if isinstance(user, UserModel) else {},
4156 '__metadata__': metadata,
4157 '__request__': request,
4158 '__model__': model,
4159 }
4161 if filter_functions:
4162 outlet_result, _ = await process_filter_functions(
4163 request=request,
4164 filter_context=None,
4165 filter_functions=filter_functions,
4166 filter_type='outlet',
4167 form_data=outlet_data,
4168 extra_params=extra_params,
4169 )
4170 else:
4171 outlet_result = outlet_data
4173 if outlet_result and outlet_result.get('messages'):
4174 if not is_unsaved_chat and messages_map:
4175 for message in outlet_result['messages']:
4176 outlet_message_id = message.get('id')
4177 if outlet_message_id and outlet_message_id in messages_map:
4178 original_message = messages_map[outlet_message_id]
4179 original_content = original_message.get('content') or get_output_text(
4180 original_message.get('output')
4181 )
4182 message_content = message.get('content') or get_output_text(message.get('output'))
4183 content_changed = original_content != message_content
4184 output_changed = message.get('output') and message.get('output') != original_message.get(
4185 'output'
4186 )
4187 if content_changed or output_changed:
4188 message_update = {
4189 'originalContent': original_content,
4190 **({'output': message['output']} if output_changed else {}),
4191 }
4192 if content_changed:
4193 message_update['content'] = message_content or ''
4194 await Chats.upsert_message_to_chat_by_id_and_message_id(
4195 chat_id,
4196 outlet_message_id,
4197 message_update,
4198 )
4200 if event_emitter:
4201 await event_emitter(
4202 {
4203 'type': 'chat:outlet',
4204 'data': {'messages': outlet_result['messages']},
4205 }
4206 )
4207 except Exception as e:
4208 log.debug('Error running outlet filters: %s', e)
4211async def non_streaming_chat_response_handler(response, ctx):
4212 request = ctx['request']
4214 user = ctx['user']
4215 metadata = ctx['metadata']
4216 events = ctx['events']
4218 event_emitter = ctx['event_emitter']
4220 response, response_data = get_response_data(response)
4221 if response_data is None:
4222 return response
4224 chat_id = metadata.get('chat_id') or ''
4225 save_to_chat = is_saved_chat_id(chat_id)
4226 continuing = bool(metadata.get('assistant_message_id'))
4228 if event_emitter:
4229 try:
4230 if 'error' in response_data:
4231 error = response_data.get('error')
4233 if isinstance(error, dict):
4234 error = error.get('detail', error)
4235 else:
4236 error = str(error)
4238 log.error('Provider returned error (non-streaming): %s', error)
4240 if save_to_chat:
4241 await Chats.upsert_message_to_chat_by_id_and_message_id(
4242 metadata['chat_id'],
4243 metadata['message_id'],
4244 {
4245 'error': {'content': error},
4246 'done': True,
4247 },
4248 )
4249 if isinstance(error, str) or isinstance(error, dict):
4250 await event_emitter(
4251 {
4252 'type': 'chat:message:error',
4253 'data': {'error': {'content': error}, 'done': True},
4254 }
4255 )
4257 if 'selected_model_id' in response_data and save_to_chat:
4258 await Chats.upsert_message_to_chat_by_id_and_message_id(
4259 metadata['chat_id'],
4260 metadata['message_id'],
4261 {
4262 'selectedModelId': response_data['selected_model_id'],
4263 },
4264 touch=False,
4265 )
4267 choices = response_data.get('choices', [])
4268 response_output = response_data.get('output')
4269 content = choices[0].get('message', {}).get('content') if choices else ''
4271 if (continuing and 'error' not in response_data) or (
4272 not continuing and choices and (content or response_output)
4273 ):
4274 if content or response_output or continuing:
4275 if not continuing:
4276 await event_emitter(
4277 {
4278 'type': 'chat:completion',
4279 'data': response_data,
4280 }
4281 )
4283 title = await Chats.get_chat_title_by_id(metadata['chat_id']) if save_to_chat else ''
4285 # Use output from backend if provided (OR-compliant backends),
4286 # otherwise generate from response content
4287 if not response_output:
4288 choice_message = choices[0].get('message', {}) if choices else {}
4289 reasoning_content = choice_message.get('reasoning_content') or choice_message.get('reasoning')
4290 reasoning_details = get_reasoning_details(choice_message)
4291 response_output = []
4292 if reasoning_content or reasoning_details:
4293 reasoning_item = {
4294 'type': 'reasoning',
4295 'id': output_id('r'),
4296 'status': 'completed',
4297 'start_tag': '<think>',
4298 'end_tag': '</think>',
4299 'attributes': {'type': 'reasoning_content'},
4300 'content': (
4301 [{'type': 'output_text', 'text': reasoning_content}] if reasoning_content else []
4302 ),
4303 'summary': None,
4304 }
4305 if reasoning_details:
4306 reasoning_item['reasoning_details'] = (
4307 reasoning_details if isinstance(reasoning_details, list) else [reasoning_details]
4308 )
4309 response_output.append(reasoning_item)
4310 response_output.append(
4311 {
4312 'type': 'message',
4313 'id': output_id('msg'),
4314 'status': 'completed',
4315 'role': 'assistant',
4316 'content': [{'type': 'output_text', 'text': content}],
4317 }
4318 )
4320 if continuing:
4321 message = ctx.get('assistant_message') or {}
4322 previous = message.get('output') or []
4323 if not previous and message.get('content'):
4324 previous = [
4325 {
4326 'type': 'message',
4327 'id': message.get('id') or output_id('msg'),
4328 'role': 'assistant',
4329 'content': [{'type': 'output_text', 'text': message['content']}],
4330 }
4331 ]
4332 response_output = list(response_output)
4333 if (
4334 previous
4335 and response_output
4336 and previous[-1].get('type') == response_output[0].get('type') == 'message'
4337 and previous[-1].get('role', 'assistant') == 'assistant'
4338 and response_output[0].get('role', 'assistant') == 'assistant'
4339 ):
4340 response_output[0] = {
4341 **previous[-1],
4342 **response_output[0],
4343 'id': previous[-1].get('id'),
4344 'content': [*previous[-1].get('content', []), *response_output[0].get('content', [])],
4345 }
4346 previous = previous[:-1]
4347 response_output = previous + response_output
4348 content = get_output_text(response_output)
4350 await event_emitter(
4351 {
4352 'type': 'chat:completion',
4353 'data': {
4354 **(response_data if continuing else {}),
4355 'done': True,
4356 'output': response_output,
4357 'title': title,
4358 },
4359 }
4360 )
4362 # Save message in the database
4363 usage = normalize_usage(response_data.get('usage', {}) or {})
4365 if save_to_chat:
4366 await Chats.upsert_message_to_chat_by_id_and_message_id(
4367 metadata['chat_id'],
4368 metadata['message_id'],
4369 {
4370 'done': True,
4371 'role': 'assistant',
4372 'output': response_output,
4373 **({'usage': usage} if usage else {}),
4374 },
4375 )
4377 await publish_chat_finished_event(request, user, metadata, title, content, response_output)
4379 ctx['assistant_message'] = {
4380 'content': content,
4381 'output': response_output,
4382 **({'usage': usage} if usage else {}),
4383 }
4384 await outlet_filter_handler(ctx)
4385 await background_tasks_handler(ctx)
4387 response = build_response_object(response, merge_events_into_response(response_data, events))
4388 except Exception as e:
4389 log.debug('Error occurred while processing request: %s', e)
4390 chat_id = metadata.get('chat_id')
4391 if getattr(request.state, 'internal', False) is not True and chat_id and is_saved_chat_id(chat_id):
4392 webui_url = await Config.get('webui.url')
4393 await publish_event(
4394 request,
4395 EVENTS.CHAT_FAILED,
4396 actor=user,
4397 subject_id=chat_id,
4398 subject_type='chat',
4399 data={
4400 'user_id': user.id,
4401 'chat_id': chat_id,
4402 'message_id': metadata.get('message_id'),
4403 'model_id': metadata.get('model_id'),
4404 'url': f'{webui_url}/c/{chat_id}' if webui_url else f'/c/{chat_id}',
4405 'message': str(e),
4406 },
4407 message='Chat failed',
4408 )
4409 pass
4411 return response
4413 choices = response_data.get('choices', [])
4414 output = response_data.get('output')
4415 content = choices[0].get('message', {}).get('content') if choices else ''
4416 if ENABLE_API_OUTLET_FILTERS and (content or output):
4417 usage = normalize_usage(response_data.get('usage', {}) or {})
4418 ctx['assistant_message'] = {
4419 **({'content': content} if content else {}),
4420 **({'output': output} if output else {}),
4421 **({'usage': usage} if usage else {}),
4422 }
4423 await outlet_filter_handler(ctx)
4425 if isinstance(response, dict):
4426 response = merge_events_into_response(response_data, events)
4428 return response
4431async def streaming_chat_response_handler(response, ctx):
4432 request = ctx['request']
4434 form_data = ctx['form_data']
4436 user = ctx['user']
4437 model = ctx['model']
4439 metadata = ctx['metadata']
4440 events = ctx['events']
4442 event_emitter = ctx['event_emitter']
4443 event_caller = ctx['event_caller']
4444 chat_id = metadata.get('chat_id') or ''
4445 save_to_chat = is_saved_chat_id(chat_id)
4446 continuing = bool(metadata.get('assistant_message_id'))
4448 extra_params = {
4449 '__event_emitter__': event_emitter,
4450 '__event_call__': event_caller,
4451 '__user__': user.model_dump() if isinstance(user, UserModel) else {},
4452 '__metadata__': metadata,
4453 '__oauth_token__': await get_system_oauth_token(request, user),
4454 '__request__': request,
4455 '__model__': model,
4456 '__chat_id__': metadata.get('chat_id'),
4457 '__message_id__': metadata.get('message_id'),
4458 }
4460 filter_functions = (
4461 await get_filter_functions(request, model, metadata.get('filter_ids', [])) if ENABLE_PLUGINS else []
4462 )
4464 # Standard streaming response handler
4465 # event_caller is optional — only needed for direct (client-side) tools
4466 # and pyodide code interpreter. Server-side tools work without it.
4467 if event_emitter:
4468 task_id = str(uuid4()) # Create a unique task ID.
4469 model_id = form_data.get('model', '')
4471 # Handle as a background task
4472 async def response_handler(response, events):
4473 filter_context = FilterContext()
4474 tag_scan_positions = {}
4475 tag_boundary_positions = {}
4476 response_stream_task_id = metadata.get('task_id') or metadata.get('message_id')
4478 def tag_output_handler(content_type, tags, output):
4479 """
4480 Detect special tags (reasoning, solution, code_interpreter) in streaming
4481 content and create corresponding OR-aligned output items directly.
4482 Operates on output items instead of content_blocks.
4484 Uses the text from the output items themselves for tag detection,
4485 eliminating state divergence between accumulated content and items.
4487 Mutates output in place; returns the rewritten output (or None
4488 if no tag was consumed) and whether a block ended.
4489 """
4490 end_flag = False
4492 def extract_attributes(tag_content):
4493 """Extract attributes from a tag if they exist."""
4494 attributes = {}
4495 if not tag_content:
4496 return attributes
4497 matches = re.findall(r'(\w+)\s*=\s*"([^"]+)"', tag_content)
4498 for key, value in matches:
4499 attributes[key] = value
4500 return attributes
4502 def get_last_text(out):
4503 """Get text from last message item, or empty string."""
4504 if out and out[-1].get('type') == 'message':
4505 parts = out[-1].get('content', [])
4506 if parts and parts[-1].get('type') == 'output_text':
4507 return parts[-1].get('text', '')
4508 return ''
4510 def set_last_text(out, text):
4511 """Set text on last message item's output_text."""
4512 if out and out[-1].get('type') == 'message':
4513 parts = out[-1].get('content', [])
4514 if parts and parts[-1].get('type') == 'output_text':
4515 parts[-1]['text'] = text
4517 def get_scanned_length(item, text):
4518 item_id = item.get('id')
4519 if not item_id:
4520 return 0
4522 scanned_length = tag_scan_positions.get((item_id, content_type), 0)
4523 return scanned_length if scanned_length <= len(text) else 0
4525 def save_scanned_length(item, text):
4526 item_id = item.get('id')
4527 if item_id:
4528 tag_scan_positions[(item_id, content_type)] = len(text)
4530 def clear_scanned_length(item):
4531 item_id = item.get('id')
4532 if item_id:
4533 tag_scan_positions.pop((item_id, content_type), None)
4534 tag_boundary_positions.pop((item_id, content_type), None)
4536 def get_tag_boundaries(item, text, scanned_length):
4537 """Index of the last '<', and of the last '>' or newline, before scanned_length."""
4538 key = (item.get('id'), content_type)
4539 scanned, last_open, last_boundary = tag_boundary_positions.get(key, (0, -1, -1))
4540 if scanned > scanned_length: # the item was rewritten, so the cached positions are stale
4541 scanned, last_open, last_boundary = 0, -1, -1
4543 if scanned < scanned_length:
4544 # only text added since the last call can move either position
4545 open_tag = text.rfind('<', scanned, scanned_length)
4546 if open_tag != -1:
4547 last_open = open_tag
4548 boundary = max(
4549 text.rfind('>', scanned, scanned_length),
4550 text.rfind('\n', scanned, scanned_length),
4551 )
4552 if boundary != -1:
4553 last_boundary = boundary
4554 tag_boundary_positions[key] = (scanned_length, last_open, last_boundary)
4556 return last_open, last_boundary
4558 # Map content_type to output item type
4559 output_type_map = {
4560 'reasoning': 'reasoning',
4561 'solution': 'message', # solution tags just produce text
4562 'code_interpreter': 'open_webui:code_interpreter',
4563 }
4564 output_item_type = output_type_map.get(content_type, content_type)
4566 last_type = output[-1].get('type', '') if output else ''
4568 if last_type == 'message':
4569 # Use the output item's own text for tag detection
4570 item = output[-1]
4571 item_text = get_last_text(output)
4572 scanned_length = get_scanned_length(item, item_text)
4573 max_start_tag_length = max((len(start_tag) for start_tag, _ in tags), default=1)
4574 search_start = max(0, scanned_length - max_start_tag_length + 1)
4576 if scanned_length and any(
4577 start_tag.startswith('<') and start_tag.endswith('>') for start_tag, _ in tags
4578 ):
4579 open_tag_start, last_tag_boundary = get_tag_boundaries(item, item_text, scanned_length)
4580 if open_tag_start > last_tag_boundary:
4581 search_start = min(search_start, open_tag_start)
4583 for start_tag, end_tag in tags:
4584 match = re.compile(_start_tag_pattern(start_tag)).search(item_text, search_start)
4585 if match:
4586 clear_scanned_length(item)
4587 try:
4588 attr_content = match.group(1) if match.group(1) else ''
4589 except Exception:
4590 attr_content = ''
4592 attributes = extract_attributes(attr_content)
4594 before_tag = item_text[: match.start()]
4595 after_tag = item_text[match.end() :]
4597 # Keep only text before the tag in the message
4598 set_last_text(output, before_tag)
4600 if not before_tag.strip():
4601 # Remove empty message item
4602 if output and output[-1].get('type') == 'message':
4603 output.pop()
4605 # Append the new output item
4606 if output_item_type == 'reasoning':
4607 output.append(
4608 {
4609 'type': 'reasoning',
4610 'id': output_id('r'),
4611 'status': 'in_progress',
4612 'start_tag': start_tag,
4613 'end_tag': end_tag,
4614 'attributes': attributes,
4615 'content': [],
4616 'summary': None,
4617 'started_at': time.time(),
4618 }
4619 )
4620 elif output_item_type == 'open_webui:code_interpreter':
4621 output.append(
4622 {
4623 'type': 'open_webui:code_interpreter',
4624 'id': output_id('ci'),
4625 'status': 'in_progress',
4626 'start_tag': start_tag,
4627 'end_tag': end_tag,
4628 'attributes': attributes,
4629 'lang': attributes.get('lang', 'python'),
4630 'code': '',
4631 'output': None,
4632 'started_at': time.time(),
4633 }
4634 )
4635 else:
4636 # solution or other text-producing tag
4637 output.append(
4638 {
4639 'type': 'message',
4640 'id': output_id('msg'),
4641 'status': 'in_progress',
4642 'role': 'assistant',
4643 'content': [{'type': 'output_text', 'text': ''}],
4644 '_tag_type': content_type,
4645 'start_tag': start_tag,
4646 'end_tag': end_tag,
4647 'attributes': attributes,
4648 'started_at': time.time(),
4649 }
4650 )
4652 if after_tag:
4653 # Set the after_tag content on the new item
4654 if output_item_type == 'reasoning':
4655 output[-1]['content'] = [{'type': 'output_text', 'text': after_tag}]
4656 elif output_item_type == 'open_webui:code_interpreter':
4657 output[-1]['code'] = after_tag
4658 else:
4659 set_last_text(output, after_tag)
4661 _, recursive_end = tag_output_handler(content_type, tags, output)
4662 if recursive_end:
4663 end_flag = True
4665 return output, end_flag
4666 else:
4667 save_scanned_length(item, item_text)
4669 elif (
4670 (last_type == 'reasoning' and content_type == 'reasoning')
4671 or (last_type == 'open_webui:code_interpreter' and content_type == 'code_interpreter')
4672 or (last_type == 'message' and output[-1].get('_tag_type') == content_type)
4673 ):
4674 item = output[-1]
4675 start_tag = item.get('start_tag', '')
4676 end_tag = item.get('end_tag', '')
4678 # Get the block content from the item itself
4679 if last_type == 'reasoning':
4680 parts = item.get('content', [])
4681 block_content = ''
4682 if parts and parts[-1].get('type') == 'output_text':
4683 block_content = parts[-1].get('text', '')
4684 elif last_type == 'open_webui:code_interpreter':
4685 block_content = item.get('code', '')
4686 else:
4687 block_content = get_last_text(output)
4689 scanned_length = get_scanned_length(item, block_content)
4690 end_tag_search_start = max(0, scanned_length - max(len(end_tag), 1) + 1)
4692 if block_content.find(end_tag, end_tag_search_start) != -1:
4693 clear_scanned_length(item)
4694 end_flag = True
4696 # Strip start and end tags from content
4697 start_tag_pattern = _start_tag_pattern(start_tag)
4698 block_content = re.sub(start_tag_pattern, '', block_content).strip()
4700 end_tag_pattern = rf'{re.escape(end_tag)}'
4701 end_tag_regex = re.compile(end_tag_pattern, re.DOTALL)
4702 split_content = end_tag_regex.split(block_content, maxsplit=1)
4704 block_content = split_content[0].strip() if split_content else ''
4705 leftover_content = split_content[1].strip() if len(split_content) > 1 else ''
4707 if block_content:
4708 # Update the item with final content
4709 if last_type == 'reasoning':
4710 item['content'] = [{'type': 'output_text', 'text': block_content}]
4711 item['ended_at'] = time.time()
4712 item['duration'] = int(item['ended_at'] - item['started_at'])
4713 item['status'] = 'completed'
4714 elif last_type == 'open_webui:code_interpreter':
4715 item['code'] = block_content
4716 item['ended_at'] = time.time()
4717 item['duration'] = int(item['ended_at'] - item['started_at'])
4718 else:
4719 set_last_text(output, block_content)
4720 item['ended_at'] = time.time()
4722 # Reset by appending a new message item for leftover
4723 output.append(
4724 {
4725 'type': 'message',
4726 'id': output_id('msg'),
4727 'status': 'in_progress',
4728 'role': 'assistant',
4729 'content': [
4730 {
4731 'type': 'output_text',
4732 'text': leftover_content,
4733 }
4734 ],
4735 }
4736 )
4737 else:
4738 # Remove the block if content is empty
4739 output.pop()
4740 output.append(
4741 {
4742 'type': 'message',
4743 'id': output_id('msg'),
4744 'status': 'in_progress',
4745 'role': 'assistant',
4746 'content': [
4747 {
4748 'type': 'output_text',
4749 'text': leftover_content,
4750 }
4751 ],
4752 }
4753 )
4754 return output, end_flag
4755 else:
4756 save_scanned_length(item, block_content)
4758 return None, False
4760 message = (
4761 ctx.get('assistant_message')
4762 if continuing
4763 else (
4764 await Chats.get_message_by_id_and_message_id(metadata['chat_id'], metadata['message_id'])
4765 if save_to_chat
4766 else None
4767 )
4768 )
4770 tool_calls = []
4772 last_assistant_message = None
4773 try:
4774 if form_data['messages'][-1]['role'] == 'assistant':
4775 last_assistant_message = get_last_assistant_message(form_data['messages'])
4776 except Exception as e:
4777 pass
4779 initial_content = (
4780 message.get('content', '') if message else last_assistant_message if last_assistant_message else ''
4781 )
4782 content_parts = [initial_content] if initial_content else []
4784 # Initialize output: use existing from message if continuing, else create new
4785 existing_output = message.get('output') if message else None
4786 prior_output = []
4787 if existing_output and metadata.get('assistant_message_id'):
4788 prior_output = list(existing_output)
4789 if (
4790 prior_output
4791 and prior_output[-1].get('type') == 'message'
4792 and prior_output[-1].get('status') == 'in_progress'
4793 ):
4794 msg_parts = prior_output[-1].get('content', [])
4795 if not msg_parts or (len(msg_parts) == 1 and not msg_parts[0].get('text', '')):
4796 prior_output.pop()
4797 output = []
4798 content_parts = []
4799 elif existing_output:
4800 output = existing_output
4801 else:
4802 # Only create an initial message item if there is content to initialize with
4803 if initial_content:
4804 output = [
4805 {
4806 'type': 'message',
4807 'id': output_id('msg'),
4808 'status': 'in_progress',
4809 'role': 'assistant',
4810 'content': [{'type': 'output_text', 'text': initial_content}],
4811 }
4812 ]
4813 else:
4814 output = []
4816 if continuing and not prior_output:
4817 # Keep the prefix separate from provider output indices and tool follow-up messages.
4818 prior_output, output = output, []
4819 content_parts = []
4821 usage = None
4822 last_response_id = None
4824 def full_output():
4825 if (
4826 continuing
4827 and prior_output
4828 and output
4829 and prior_output[-1].get('type') == output[0].get('type') == 'message'
4830 and prior_output[-1].get('role', 'assistant') == 'assistant'
4831 and output[0].get('role', 'assistant') == 'assistant'
4832 ):
4833 return [
4834 *prior_output[:-1],
4835 {
4836 **prior_output[-1],
4837 **output[0],
4838 'id': prior_output[-1].get('id'),
4839 'content': [*prior_output[-1].get('content', []), *output[0].get('content', [])],
4840 },
4841 *output[1:],
4842 ]
4843 return prior_output + output if prior_output else output
4845 def get_message_error_content(error):
4846 if isinstance(error, HTTPException):
4847 error = error.detail
4848 elif isinstance(error, dict):
4849 error = error.get('detail', error)
4850 else:
4851 error = str(error)
4853 return error if isinstance(error, (str, dict)) else str(error)
4855 async def emit_message_error(error_content):
4856 if save_to_chat:
4857 await Chats.upsert_message_to_chat_by_id_and_message_id(
4858 metadata['chat_id'],
4859 metadata['message_id'],
4860 {'error': {'content': error_content}},
4861 )
4862 await event_emitter(
4863 {
4864 'type': 'chat:message:error',
4865 'data': {'error': {'content': error_content}},
4866 }
4867 )
4869 reasoning_tags_param = metadata.get('params', {}).get('reasoning_tags')
4870 DETECT_REASONING_TAGS = reasoning_tags_param is not False
4872 # Legacy tool-calling only: native FC gets execute_code as a builtin tool.
4873 # Same five authz gates as utils/tools.py get_builtin_tools.
4874 features = metadata.get('features', {}) or {}
4875 model_capabilities = model.get('info', {}).get('meta', {}).get('capabilities') or {}
4876 builtin_tools_meta = model.get('info', {}).get('meta', {}).get('builtinTools', {})
4877 DETECT_CODE_INTERPRETER = (
4878 metadata.get('params', {}).get('function_calling') == 'legacy'
4879 and bool(features.get('code_interpreter'))
4880 and builtin_tools_meta.get('code_interpreter', True)
4881 and await Config.get('code_interpreter.enable')
4882 and model_capabilities.get('code_interpreter', True)
4883 and (
4884 getattr(user, 'role', None) == 'admin'
4885 or await has_permission(
4886 getattr(user, 'id', ''),
4887 'features.code_interpreter',
4888 await Config.get('user.permissions'),
4889 )
4890 )
4891 )
4893 reasoning_tags = []
4894 if DETECT_REASONING_TAGS:
4895 if isinstance(reasoning_tags_param, list) and len(reasoning_tags_param) == 2:
4896 reasoning_tags = [(reasoning_tags_param[0], reasoning_tags_param[1])]
4897 else:
4898 reasoning_tags = DEFAULT_REASONING_TAGS
4900 try:
4901 for event in events:
4902 await event_emitter(
4903 {
4904 'type': 'chat:completion',
4905 'data': event,
4906 }
4907 )
4909 # Save message in the database
4910 if save_to_chat:
4911 await Chats.upsert_message_to_chat_by_id_and_message_id(
4912 metadata['chat_id'],
4913 metadata['message_id'],
4914 {
4915 **event,
4916 },
4917 )
4919 async def stream_body_handler(response, form_data):
4920 nonlocal usage
4921 nonlocal output
4922 nonlocal prior_output
4923 nonlocal last_response_id
4925 response_tool_calls = []
4927 delta_count = 0
4928 delta_chunk_size = max(
4929 CHAT_RESPONSE_STREAM_DELTA_CHUNK_SIZE,
4930 int(metadata.get('params', {}).get('stream_delta_chunk_size') or 1),
4931 )
4932 last_delta_data = None
4933 last_delta_type = None
4934 last_delta_key = None
4936 joined_content = ''
4937 joined_part_count = 0
4939 async def save_current_response_stream(stream_output: list | None = None):
4940 nonlocal joined_content
4941 nonlocal joined_part_count
4943 if not chat_id or not metadata.get('message_id'):
4944 return
4946 # content_parts is append-only, so its length tells us when the join is stale
4947 if joined_part_count != len(content_parts):
4948 joined_content = ''.join(content_parts)
4949 joined_part_count = len(content_parts)
4951 current_stream_output = stream_output if stream_output is not None else full_output()
4952 await save_response_stream(
4953 request.app.state.redis,
4954 response_stream_task_id,
4955 chat_id,
4956 metadata.get('message_id'),
4957 get_output_text(current_stream_output)
4958 if continuing
4959 else joined_content or get_output_text(current_stream_output),
4960 current_stream_output,
4961 )
4963 def get_response_delta_key(delta_data: dict):
4964 event_type = delta_data.get('type', '')
4965 if not event_type.startswith('response.') or not event_type.endswith('.delta'):
4966 return None
4967 return (
4968 event_type,
4969 delta_data.get('item_id'),
4970 delta_data.get('output_index'),
4971 delta_data.get('content_index'),
4972 delta_data.get('summary_index'),
4973 )
4975 def get_response_data_with_full_output_index(response_data: dict):
4976 if prior_output and isinstance(response_data.get('output_index'), int):
4977 return {
4978 **response_data,
4979 'output_index': response_data['output_index'] + len(prior_output),
4980 }
4981 if prior_output and isinstance(response_data.get('response'), dict):
4982 # response.output is this round only; the response.completed reducer drops earlier rounds
4983 return {
4984 **response_data,
4985 'response': {
4986 key: value for key, value in response_data['response'].items() if key != 'output'
4987 },
4988 }
4989 return response_data
4991 async def flush_pending_delta_data(threshold: int = 0):
4992 nonlocal delta_count
4993 nonlocal last_delta_data
4994 nonlocal last_delta_type
4995 nonlocal last_delta_key
4997 if delta_count >= threshold and last_delta_data:
4998 await event_emitter(
4999 {
5000 'type': 'chat:completion' if continuing else 'response:completion',
5001 'data': {'output': full_output(), 'type': last_delta_data.get('type')}
5002 if continuing
5003 else last_delta_data,
5004 }
5005 )
5006 await save_current_response_stream()
5007 delta_count = 0
5008 last_delta_data = None
5009 last_delta_type = None
5010 last_delta_key = None
5012 async def queue_pending_delta_data(delta_data: dict, delta_type: str):
5013 nonlocal delta_count
5014 nonlocal last_delta_data
5015 nonlocal last_delta_type
5016 nonlocal last_delta_key
5018 delta_data = get_response_data_with_full_output_index(delta_data)
5019 delta_key = get_response_delta_key(delta_data)
5020 if (
5021 last_delta_data
5022 and last_delta_key == delta_key
5023 and isinstance(last_delta_data.get('delta'), str)
5024 and isinstance(delta_data.get('delta'), str)
5025 ):
5026 append_to_text_field(last_delta_data, 'delta', delta_data['delta'])
5027 delta_count += 1
5028 else:
5029 if last_delta_data and (last_delta_type != delta_type or last_delta_key != delta_key):
5030 await flush_pending_delta_data()
5032 delta_count += 1
5033 last_delta_data = delta_data
5034 last_delta_type = delta_type
5035 last_delta_key = delta_key
5037 if delta_count >= delta_chunk_size:
5038 await flush_pending_delta_data(delta_chunk_size)
5040 async def emit_response_completion_event(response_data: dict, stream_output: list | None = None):
5041 if response_data.get('type', '').endswith('.delta'):
5042 await queue_pending_delta_data(
5043 response_data,
5044 response_data.get('type', 'response.delta'),
5045 )
5046 return
5048 response_data = get_response_data_with_full_output_index(response_data)
5049 await flush_pending_delta_data()
5050 await event_emitter(
5051 {
5052 'type': 'chat:completion' if continuing else 'response:completion',
5053 'data': {'output': full_output()}
5054 if continuing
5055 else get_response_completion_event_data(response_data),
5056 }
5057 )
5058 await save_current_response_stream(stream_output)
5060 filter_extra_params = {'__body__': form_data, **extra_params} if filter_functions else None
5062 async for line in response.body_iterator:
5063 line = line.decode('utf-8', 'replace') if isinstance(line, bytes) else line
5064 data = line
5066 # Skip empty lines
5067 if not data or data.isspace():
5068 continue
5070 # "data:" is the prefix for each event
5071 if not data.startswith('data:'):
5072 # Some upstreams return plain JSON error lines in a streaming response
5073 # (without SSE `data:` prefix). Try to normalize these into standard
5074 # error events so frontend and DB paths still receive them.
5075 try:
5076 raw_obj = JSONCodec.loads(data)
5077 raw_error = raw_obj.get('error') if isinstance(raw_obj, dict) else None
5078 if raw_error:
5079 if save_to_chat:
5080 try:
5081 await Chats.upsert_message_to_chat_by_id_and_message_id(
5082 metadata['chat_id'],
5083 metadata['message_id'],
5084 {
5085 'error': {'content': raw_error},
5086 },
5087 )
5088 except Exception:
5089 pass
5090 await event_emitter({'type': 'chat:completion', 'data': {'error': raw_error}})
5091 except Exception:
5092 pass
5093 continue
5095 # Remove the "data:" prefix
5096 data = data[5:].strip()
5098 try:
5099 data = JSONCodec.loads(data)
5101 if filter_functions:
5102 data, _ = await process_filter_functions(
5103 request=request,
5104 filter_context=filter_context,
5105 filter_functions=filter_functions,
5106 filter_type='stream',
5107 form_data=data,
5108 extra_params=filter_extra_params,
5109 )
5111 if data:
5112 if 'event' in data and not getattr(request.state, 'direct', False):
5113 await event_emitter(data.get('event', {}))
5115 if 'selected_model_id' in data:
5116 model_id = data['selected_model_id']
5117 if save_to_chat:
5118 await Chats.upsert_message_to_chat_by_id_and_message_id(
5119 metadata['chat_id'],
5120 metadata['message_id'],
5121 {
5122 'selectedModelId': model_id,
5123 },
5124 touch=False,
5125 )
5126 await event_emitter(
5127 {
5128 'type': 'chat:completion',
5129 'data': data,
5130 }
5131 )
5132 # Check for Responses API events (type field starts with "response.")
5133 elif data.get('type', '').startswith('response.'):
5134 response_data_type = data.get('type', '')
5135 response_data_is_delta = response_data_type.endswith('.delta')
5136 output, response_metadata = handle_responses_streaming_event(data, output)
5138 if not response_data_is_delta:
5139 await flush_pending_delta_data()
5141 # Emit citation sources from finalized output items
5142 # (mirrors Chat Completions annotation handling at delta level)
5143 if response_data_type == 'response.output_item.done':
5144 item = data.get('item', {})
5145 if item.get('type') == 'message':
5146 for part in item.get('content', []):
5147 for annotation in part.get('annotations', []):
5148 if annotation.get('type') == 'url_citation':
5149 # Handle both flat (Responses API) and nested (Chat Completions) formats
5150 url_citation = annotation.get('url_citation', annotation)
5152 url = url_citation.get('url', '')
5153 title = url_citation.get('title', url)
5155 if url:
5156 await event_emitter(
5157 {
5158 'type': 'source',
5159 'data': {
5160 'source': {
5161 'name': title,
5162 'url': url,
5163 },
5164 'document': [title],
5165 'metadata': [
5166 {
5167 'source': url,
5168 'name': title,
5169 }
5170 ],
5171 },
5172 }
5173 )
5175 # Merge any metadata (usage, etc.)
5176 # Strip 'done' — response.completed emits
5177 # it but we may still need to execute tool
5178 # calls. The outer middleware manages the
5179 # actual completion signal.
5180 if response_metadata:
5181 if ENABLE_RESPONSES_API_STATEFUL:
5182 response_id = response_metadata.pop('response_id', None)
5183 if response_id:
5184 last_response_id = response_id
5186 # Normalize and capture usage for DB persistence
5187 if response_metadata.get('usage'):
5188 usage = merge_usage(usage, response_metadata['usage'])
5189 response_metadata['usage'] = usage
5191 if response_metadata.get('error'):
5192 await event_emitter(
5193 {
5194 'type': 'chat:completion',
5195 'data': {'error': response_metadata['error']},
5196 }
5197 )
5199 await emit_response_completion_event(data)
5201 if response_metadata and response_metadata.get('usage'):
5202 await event_emitter(
5203 {
5204 'type': 'chat:completion',
5205 'data': {'usage': usage},
5206 }
5207 )
5208 continue
5209 else:
5210 choices = data.get('choices', [])
5212 # Normalize usage data to standard format
5213 raw_usage = data.get('usage', {}) or {}
5214 raw_usage.update(data.get('timings', {})) # llama.cpp
5215 if raw_usage:
5216 usage = merge_usage(usage, raw_usage)
5217 await event_emitter(
5218 {
5219 'type': 'chat:completion',
5220 'data': {
5221 'usage': usage,
5222 },
5223 }
5224 )
5226 if not choices:
5227 error = data.get('error', {})
5228 if error:
5229 log.error('Provider returned error (streaming): %s', error)
5230 if save_to_chat:
5231 try:
5232 await Chats.upsert_message_to_chat_by_id_and_message_id(
5233 metadata['chat_id'],
5234 metadata['message_id'],
5235 {
5236 'error': {'content': error},
5237 },
5238 )
5239 except Exception:
5240 pass
5241 await event_emitter(
5242 {
5243 'type': 'chat:completion',
5244 'data': {
5245 'error': error,
5246 },
5247 }
5248 )
5249 continue
5251 delta = choices[0].get('delta', {})
5252 delta_type = 'content'
5254 # Handle delta annotations
5255 annotations = delta.get('annotations')
5256 if annotations:
5257 for annotation in annotations:
5258 if (
5259 annotation.get('type') == 'url_citation'
5260 and 'url_citation' in annotation
5261 ):
5262 url_citation = annotation['url_citation']
5264 url = url_citation.get('url', '')
5265 title = url_citation.get('title', url)
5267 await event_emitter(
5268 {
5269 'type': 'source',
5270 'data': {
5271 'source': {
5272 'name': title,
5273 'url': url,
5274 },
5275 'document': [title],
5276 'metadata': [
5277 {
5278 'source': url,
5279 'name': title,
5280 }
5281 ],
5282 },
5283 }
5284 )
5286 delta_tool_calls = delta.get('tool_calls', None)
5287 if delta_tool_calls:
5288 for delta_tool_call in delta_tool_calls:
5289 tool_call_index = delta_tool_call.get('index')
5291 if tool_call_index is not None:
5292 # Check if the tool call already exists
5293 current_response_tool_call = None
5294 for response_tool_call in response_tool_calls:
5295 if response_tool_call.get('index') == tool_call_index:
5296 current_response_tool_call = response_tool_call
5297 break
5299 if current_response_tool_call is None:
5300 # Add the new tool call
5301 delta_tool_call.setdefault('function', {})
5302 if delta_tool_call['function'].get('name') is None:
5303 delta_tool_call['function']['name'] = ''
5304 delta_tool_call['id'] = delta_tool_call.get('id') or output_id('fc')
5305 delta_arguments = delta_tool_call['function'].get('arguments')
5306 if not isinstance(delta_arguments, str):
5307 delta_tool_call['function']['arguments'] = (
5308 ''
5309 if delta_arguments is None
5310 else JSONCodec.dumps(delta_arguments)
5311 )
5312 response_tool_calls.append(delta_tool_call)
5313 else:
5314 # Update the existing tool call
5315 delta_name = delta_tool_call.get('function', {}).get('name')
5316 delta_arguments = delta_tool_call.get('function', {}).get(
5317 'arguments'
5318 )
5320 if delta_name:
5321 current_response_tool_call['function']['name'] = delta_name
5323 if delta_arguments is not None:
5324 if not isinstance(delta_arguments, str):
5325 delta_arguments = JSONCodec.dumps(delta_arguments)
5326 current_response_tool_call.setdefault('function', {})
5327 if not isinstance(
5328 current_response_tool_call['function'].get('arguments'),
5329 str,
5330 ):
5331 current_response_tool_call['function']['arguments'] = ''
5332 append_to_text_field(
5333 current_response_tool_call['function'],
5334 'arguments',
5335 delta_arguments,
5336 )
5338 # Emit pending tool calls in real-time as Responses events.
5339 if response_tool_calls:
5340 output_by_call_id = {
5341 item.get('call_id'): (idx, item)
5342 for idx, item in enumerate(output)
5343 if item.get('type') == 'function_call'
5344 }
5346 for tc in response_tool_calls:
5347 call_id = tc.get('id') or output_id('fc')
5348 tc['id'] = call_id
5349 func = tc.get('function', {})
5350 if call_id in output_by_call_id:
5351 output_index, item = output_by_call_id[call_id]
5352 item['name'] = func.get('name', item.get('name', ''))
5353 item['arguments'] = func.get('arguments', item.get('arguments', ''))
5354 item['status'] = 'in_progress'
5355 else:
5356 output_index = len(output)
5357 item = {
5358 'type': 'function_call',
5359 'id': call_id,
5360 'call_id': call_id,
5361 'name': func.get('name', ''),
5362 'arguments': '',
5363 'status': 'in_progress',
5364 }
5365 output.append(item)
5366 output_by_call_id[call_id] = (output_index, item)
5367 await emit_response_completion_event(
5368 {
5369 'type': 'response.output_item.added',
5370 'output_index': output_index,
5371 'item': item.copy(),
5372 }
5373 )
5374 item['arguments'] = func.get('arguments', '')
5376 for delta_tool_call in delta_tool_calls:
5377 tool_call_index = delta_tool_call.get('index')
5378 current_response_tool_call = next(
5379 (
5380 tc
5381 for tc in response_tool_calls
5382 if tc.get('index') == tool_call_index
5383 ),
5384 None,
5385 )
5386 if not current_response_tool_call:
5387 continue
5388 call_id = current_response_tool_call.get('id')
5389 output_index, _ = output_by_call_id.get(call_id, (len(output) - 1, {}))
5390 delta_arguments = delta_tool_call.get('function', {}).get('arguments')
5391 if delta_arguments is not None:
5392 if not isinstance(delta_arguments, str):
5393 delta_arguments = JSONCodec.dumps(delta_arguments)
5394 await emit_response_completion_event(
5395 {
5396 'type': 'response.function_call_arguments.delta',
5397 'item_id': call_id,
5398 'output_index': output_index,
5399 'delta': delta_arguments,
5400 }
5401 )
5403 await save_current_response_stream()
5404 data = None
5405 delta_type = 'tool_call'
5407 delta_images = delta.get('images')
5408 image_urls = (
5409 await get_image_urls(delta_images, request, metadata, user)
5410 if delta_images
5411 else []
5412 )
5413 if image_urls:
5414 image_file_list = [{'type': 'image', 'url': url} for url in image_urls]
5415 message_files = image_file_list
5416 if save_to_chat:
5417 message_files = await Chats.add_message_files_by_id_and_message_id(
5418 metadata['chat_id'],
5419 metadata['message_id'],
5420 image_file_list,
5421 )
5422 if message_files is None:
5423 message_files = image_file_list
5425 await event_emitter(
5426 {
5427 'type': 'files',
5428 'data': {'files': message_files},
5429 }
5430 )
5432 # content and reasoning deltas are raw JSON: a stream filter can make them any type
5433 value = delta.get('content')
5434 if value and not isinstance(value, str):
5435 value = f'{value}'
5437 reasoning_content = (
5438 delta.get('reasoning_content')
5439 or delta.get('reasoning')
5440 or delta.get('thinking')
5441 )
5442 if reasoning_content and not isinstance(reasoning_content, str):
5443 reasoning_content = f'{reasoning_content}'
5444 reasoning_details = get_reasoning_details(delta)
5445 reasoning_detail_items = (
5446 [item for item in reasoning_details if isinstance(item, dict)]
5447 if isinstance(reasoning_details, list)
5448 else [reasoning_details]
5449 if isinstance(reasoning_details, dict)
5450 else []
5451 )
5452 existing_reasoning_item = next(
5453 (item for item in reversed(output) if item.get('type') == 'reasoning'),
5454 None,
5455 )
5456 message_index = next(
5457 (i for i, item in enumerate(output) if item.get('type') == 'message'),
5458 None,
5459 )
5460 if reasoning_content or (
5461 reasoning_detail_items
5462 and (
5463 existing_reasoning_item
5464 or any(
5465 item.get('text') or item.get('summary') or item.get('data')
5466 for item in reasoning_detail_items
5467 )
5468 )
5469 ):
5470 reasoning_item = (
5471 existing_reasoning_item
5472 if (reasoning_detail_items and not reasoning_content)
5473 or message_index is not None
5474 else None
5475 )
5477 if reasoning_item is None:
5478 if not output or output[-1].get('type') != 'reasoning':
5479 reasoning_item = {
5480 'type': 'reasoning',
5481 'id': output_id('r'),
5482 'status': 'in_progress',
5483 'start_tag': '<think>',
5484 'end_tag': '</think>',
5485 'attributes': {'type': 'reasoning_content'},
5486 'content': [],
5487 'summary': None,
5488 'started_at': time.time(),
5489 }
5490 if message_index is not None:
5491 reasoning_item['ended_at'] = time.time()
5492 reasoning_item['duration'] = 0
5493 reasoning_item['status'] = 'completed'
5494 output.insert(message_index, reasoning_item)
5495 else:
5496 output.append(reasoning_item)
5497 else:
5498 reasoning_item = output[-1]
5500 if reasoning_content:
5501 # Append to reasoning content
5502 parts = reasoning_item.get('content', [])
5503 if parts and parts[-1].get('type') == 'output_text':
5504 append_to_text_field(parts[-1], 'text', reasoning_content)
5505 else:
5506 reasoning_item['content'] = [
5507 {
5508 'type': 'output_text',
5509 'text': reasoning_content,
5510 }
5511 ]
5513 reasoning_index = output.index(reasoning_item)
5514 data = {
5515 'type': 'response.reasoning_text.delta',
5516 'item_id': reasoning_item.get('id'),
5517 'output_index': reasoning_index,
5518 'content_index': max(
5519 len(reasoning_item.get('content', [])) - 1,
5520 0,
5521 ),
5522 'delta': reasoning_content,
5523 }
5524 delta_type = 'response.reasoning_text.delta'
5526 if reasoning_detail_items:
5527 merge_streamed_reasoning_details(
5528 reasoning_item.setdefault('reasoning_details', []),
5529 reasoning_detail_items,
5530 )
5531 await save_current_response_stream()
5532 # Providers such as OpenRouter send reasoning_details
5533 # alongside the reasoning text: only drop the event when
5534 # the details were all there was to report, otherwise the
5535 # reasoning delta never reaches the client.
5536 if not reasoning_content:
5537 data = None
5539 if value:
5540 if (
5541 output
5542 and output[-1].get('type') == 'reasoning'
5543 and output[-1].get('attributes', {}).get('type') == 'reasoning_content'
5544 ):
5545 reasoning_item = output[-1]
5546 reasoning_item['ended_at'] = time.time()
5547 reasoning_item['duration'] = int(
5548 reasoning_item['ended_at'] - reasoning_item['started_at']
5549 )
5550 reasoning_item['status'] = 'completed'
5552 output.append(
5553 {
5554 'type': 'message',
5555 'id': output_id('msg'),
5556 'status': 'in_progress',
5557 'role': 'assistant',
5558 'content': [
5559 {
5560 'type': 'output_text',
5561 'text': '',
5562 }
5563 ],
5564 }
5565 )
5567 if ENABLE_CHAT_RESPONSE_BASE64_IMAGE_URL_CONVERSION:
5568 value = await convert_markdown_base64_images(
5569 request,
5570 value,
5571 {
5572 'chat_id': metadata.get('chat_id', None),
5573 'message_id': metadata.get('message_id', None),
5574 },
5575 user,
5576 )
5578 # closure-cell str += recopies per chunk; append + join once at read is O(n)
5579 content_parts.append(value)
5581 # Check if we're inside a tag-based block
5582 # (reasoning, code_interpreter, or solution).
5583 # If so, append to the existing in-progress
5584 # item instead of creating a new message —
5585 # otherwise tag_output_handler re-detects the
5586 # start tag on every chunk and fragments the
5587 # output.
5588 last_item = output[-1] if output else None
5589 last_item_type = last_item.get('type', '') if last_item else ''
5590 inside_tag_block = (
5591 last_item is not None
5592 and last_item.get('status') == 'in_progress'
5593 and last_item.get('attributes', {}).get('type') != 'reasoning_content'
5594 and (
5595 last_item_type == 'reasoning'
5596 or last_item_type == 'open_webui:code_interpreter'
5597 or (
5598 last_item_type == 'message'
5599 and last_item.get('_tag_type') is not None
5600 )
5601 )
5602 )
5604 if inside_tag_block:
5605 # Append to the existing tag-based item
5606 if last_item_type == 'open_webui:code_interpreter':
5607 last_item['code'] = last_item.get('code', '') + value
5608 elif last_item_type == 'reasoning':
5609 parts = last_item.get('content', [])
5610 if parts and parts[-1].get('type') == 'output_text':
5611 append_to_text_field(parts[-1], 'text', value)
5612 else:
5613 last_item['content'] = [
5614 {
5615 'type': 'output_text',
5616 'text': value,
5617 }
5618 ]
5619 else:
5620 # solution or other _tag_type message
5621 msg_parts = last_item.get('content', [])
5622 if msg_parts and msg_parts[-1].get('type') == 'output_text':
5623 append_to_text_field(msg_parts[-1], 'text', value)
5624 else:
5625 last_item['content'] = [
5626 {
5627 'type': 'output_text',
5628 'text': value,
5629 }
5630 ]
5631 else:
5632 if not output or output[-1].get('type') != 'message':
5633 output.append(
5634 {
5635 'type': 'message',
5636 'id': output_id('msg'),
5637 'status': 'in_progress',
5638 'role': 'assistant',
5639 'content': [
5640 {
5641 'type': 'output_text',
5642 'text': '',
5643 }
5644 ],
5645 }
5646 )
5648 # Append value to last message item's text
5649 msg_parts = output[-1].get('content', [])
5650 if msg_parts and msg_parts[-1].get('type') == 'output_text':
5651 append_to_text_field(msg_parts[-1], 'text', value)
5652 else:
5653 output[-1]['content'] = [
5654 {
5655 'type': 'output_text',
5656 'text': value,
5657 }
5658 ]
5660 tag_output = None
5662 if DETECT_REASONING_TAGS:
5663 tag_output, _ = tag_output_handler(
5664 'reasoning',
5665 reasoning_tags,
5666 output,
5667 )
5669 solution_output, _ = tag_output_handler(
5670 'solution',
5671 DEFAULT_SOLUTION_TAGS,
5672 output,
5673 )
5674 if solution_output is not None:
5675 tag_output = solution_output
5677 if DETECT_CODE_INTERPRETER:
5678 code_output, end = tag_output_handler(
5679 'code_interpreter',
5680 DEFAULT_CODE_INTERPRETER_TAGS,
5681 output,
5682 )
5683 if code_output is not None:
5684 tag_output = code_output
5686 if end:
5687 break
5689 target_index = len(output) - 1
5690 target_item = output[target_index] if target_index >= 0 else {}
5691 target_content = target_item.get('content', [])
5692 content_index = max(len(target_content) - 1, 0)
5693 delta_event_type = (
5694 'response.reasoning_text.delta'
5695 if target_item.get('type') == 'reasoning'
5696 else 'response.output_text.delta'
5697 )
5698 data = {
5699 'type': delta_event_type,
5700 'item_id': target_item.get('id'),
5701 'output_index': target_index,
5702 'content_index': content_index,
5703 'delta': value,
5704 }
5705 delta_type = delta_event_type
5707 # the raw chunk still carries the tag text: resend the cleaned output instead
5708 if tag_output is not None:
5709 await flush_pending_delta_data()
5710 await event_emitter(
5711 {'type': 'chat:completion', 'data': {'output': full_output()}}
5712 )
5713 await save_current_response_stream()
5714 data = None
5716 if delta and data:
5717 await queue_pending_delta_data(data, delta_type)
5718 elif data:
5719 await event_emitter(
5720 {
5721 'type': 'chat:completion',
5722 'data': data,
5723 }
5724 )
5725 except (asyncio.CancelledError, KeyboardInterrupt):
5726 raise
5727 except Exception as e:
5728 done = 'data: [DONE]' in line
5729 if done:
5730 pass
5731 else:
5732 log.debug('Error: %s', e)
5733 continue
5734 await flush_pending_delta_data()
5736 if output:
5737 # Clean up the last message item
5738 if output[-1].get('type') == 'message' and not continuing:
5739 parts = output[-1].get('content', [])
5740 if parts and parts[-1].get('type') == 'output_text':
5741 parts[-1]['text'] = parts[-1]['text'].strip()
5743 if not parts[-1]['text']:
5744 output.pop()
5746 if not output:
5747 output.append(
5748 {
5749 'type': 'message',
5750 'id': output_id('msg'),
5751 'status': 'in_progress',
5752 'role': 'assistant',
5753 'content': [{'type': 'output_text', 'text': ''}],
5754 }
5755 )
5757 if output[-1].get('type') == 'reasoning':
5758 reasoning_item = output[-1]
5759 if reasoning_item.get('ended_at') is None:
5760 reasoning_item['ended_at'] = time.time()
5761 if reasoning_item.get('started_at') is not None:
5762 reasoning_item['duration'] = int(
5763 reasoning_item['ended_at'] - reasoning_item['started_at']
5764 )
5765 reasoning_item['status'] = 'completed'
5767 if response_tool_calls:
5768 for tc in response_tool_calls:
5769 call_id = tc.get('id', '')
5770 arguments = tc.get('function', {}).get('arguments', '{}')
5771 for output_index, item in enumerate(output):
5772 if item.get('type') == 'function_call' and item.get('call_id') == call_id:
5773 item['arguments'] = arguments
5774 item['status'] = 'completed'
5775 await emit_response_completion_event(
5776 {
5777 'type': 'response.function_call_arguments.done',
5778 'item_id': item.get('id'),
5779 'output_index': output_index,
5780 'arguments': arguments,
5781 }
5782 )
5783 await emit_response_completion_event(
5784 {
5785 'type': 'response.output_item.done',
5786 'output_index': output_index,
5787 'item': item.copy(),
5788 }
5789 )
5790 break
5791 tool_calls.append(_split_tool_calls(response_tool_calls))
5793 # Responses API path: extract function_call items from output
5794 if not response_tool_calls and output:
5795 # Collect call_ids that already have results,
5796 # including those from prior_output so we don't
5797 # re-process tool calls from a previous turn.
5798 handled_call_ids = {
5799 item.get('call_id')
5800 for item in (prior_output + output)
5801 if item.get('type') == 'function_call_output'
5802 }
5803 responses_api_tool_calls = []
5804 for item in output:
5805 call_id = item.get('call_id') or item.get('id') or output_id('fc')
5806 if item.get('type') == 'function_call' and call_id not in handled_call_ids:
5807 arguments = item.get('arguments', '{}')
5808 responses_api_tool_calls.append(
5809 {
5810 'id': call_id,
5811 'index': len(responses_api_tool_calls),
5812 'function': {
5813 'name': item.get('name', ''),
5814 'arguments': (
5815 arguments if isinstance(arguments, str) else JSONCodec.dumps(arguments)
5816 ),
5817 },
5818 }
5819 )
5820 if responses_api_tool_calls:
5821 tool_calls.append(_split_tool_calls(responses_api_tool_calls))
5823 output_start = len(prior_output)
5824 try:
5825 await stream_body_handler(response, form_data)
5826 finally:
5827 if response.background:
5828 await response.background()
5830 tool_call_iterations = 0
5831 max_tool_call_iterations = getattr(
5832 request.state,
5833 'max_tool_call_iterations',
5834 CHAT_RESPONSE_MAX_TOOL_CALL_ITERATIONS,
5835 )
5836 tool_call_sources = [] # Track citation sources from tool results
5837 all_tool_call_sources = [] # Accumulated sources across all iterations
5838 user_message = get_last_user_message(form_data['messages'])
5840 # Check if citations are enabled for this model
5841 citations_enabled = (model.get('info', {}).get('meta', {}).get('capabilities') or {}).get(
5842 'citations', True
5843 )
5845 # Use the pre-RAG system content captured before the
5846 # initial file-source injection in process_chat_payload.
5847 # This ensures restore truly undoes the RAG template.
5848 original_system_content = metadata.get('system_prompt')
5849 if original_system_content is None:
5850 original_system_message = get_system_message(form_data['messages'])
5851 original_system_content = (
5852 get_content_from_message(original_system_message) if original_system_message else None
5853 )
5855 async def emit_output():
5856 # Channels publish whole messages; Continue can merge into the preceding item.
5857 snapshot = continuing or (metadata.get('chat_id') or '').startswith('channel:')
5858 frontend_output = []
5859 for item in full_output() if snapshot else full_output()[output_start:]:
5860 if item.get('type') == 'function_call_output':
5861 # input_image parts are base64 data URIs only for the LLM, via convert_output_to_messages
5862 item = {
5863 **item,
5864 'output': [
5865 part for part in item.get('output', []) if part.get('type') != 'input_image'
5866 ],
5867 }
5868 frontend_output.append(item)
5870 if snapshot:
5871 await event_emitter(
5872 {
5873 'type': 'chat:completion',
5874 'data': {'output': frontend_output, 'flush': True},
5875 }
5876 )
5877 return
5879 for output_index, item in enumerate(frontend_output, start=output_start):
5880 await event_emitter(
5881 {
5882 'type': 'response:completion',
5883 'data': {
5884 'type': 'response.output_item.done',
5885 'output_index': output_index,
5886 'item': item,
5887 },
5888 }
5889 )
5891 while tool_calls and (
5892 max_tool_call_iterations is None or tool_call_iterations < max_tool_call_iterations
5893 ):
5894 tool_call_iterations += 1
5896 response_tool_calls = tool_calls.pop(0)
5897 ask_user_staged, ask_user_error = stage_ask_user_tool_calls(response_tool_calls, output, output_id)
5898 if ask_user_error:
5899 response_tool_calls = [
5900 tool_call
5901 for tool_call in response_tool_calls
5902 if tool_call.get('function', {}).get('name') != 'ask_user'
5903 ]
5904 elif ask_user_staged:
5905 if is_saved_chat_id(metadata.get('chat_id')) and metadata.get('message_id'):
5906 await pause_for_tool_approval(
5907 metadata['chat_id'],
5908 metadata['message_id'],
5909 full_output(),
5910 form_data,
5911 metadata,
5912 )
5913 await event_emitter({'type': 'chat:completion', 'data': {'output': full_output()}})
5914 return
5916 # Append function_call items for each tool call
5917 # (Responses API already has them from streaming, so skip duplicates)
5918 existing_call_ids = {item.get('call_id') for item in output if item.get('type') == 'function_call'}
5919 for tc in response_tool_calls:
5920 call_id = tc.get('id', '')
5921 if call_id not in existing_call_ids:
5922 func = tc.get('function', {})
5923 output.append(
5924 {
5925 'type': 'function_call',
5926 'id': call_id or output_id('fc'),
5927 'call_id': call_id,
5928 'name': func.get('name', ''),
5929 'arguments': func.get('arguments', '{}'),
5930 'status': 'in_progress',
5931 }
5932 )
5934 tool_approval_mode = metadata.get('params', {}).get('tool_approval_mode', 'full')
5935 if (
5936 response_tool_calls
5937 and tool_approval_mode == 'ask'
5938 and is_saved_chat_id(metadata.get('chat_id'))
5939 and metadata.get('message_id')
5940 ):
5941 await pause_for_tool_approval(
5942 metadata['chat_id'],
5943 metadata['message_id'],
5944 full_output(),
5945 form_data,
5946 metadata,
5947 )
5948 await event_emitter(
5949 {
5950 'type': 'chat:completion',
5951 'data': {
5952 'output': full_output(),
5953 },
5954 }
5955 )
5956 return
5958 await emit_output()
5960 tools = metadata.get('tools', {})
5962 results = []
5964 def parse_tool_params(tool_call):
5965 tool_args = tool_call.get('function', {}).get('arguments', '{}')
5966 params = {}
5967 if tool_args and tool_args.strip():
5968 try:
5969 params = JSONCodec.loads(tool_args)
5970 except Exception:
5971 try:
5972 params = ast.literal_eval(tool_args)
5973 except Exception as e:
5974 log.debug(e)
5975 return None
5976 if not isinstance(params, dict):
5977 raise ValueError('Tool call arguments must be a JSON object.')
5978 tool_call.setdefault('function', {})['arguments'] = JSONCodec.dumps(params)
5979 return params
5981 async def execute_tool_call(tool_call):
5982 name = tool_call.get('function', {}).get('name', '')
5983 try:
5984 params = parse_tool_params(tool_call)
5985 except ValueError:
5986 return (
5987 {},
5988 f'Error: Tool call arguments for `{name}` must be a JSON object. Please try again.',
5989 None,
5990 None,
5991 False,
5992 )
5993 if params is None:
5994 return {}, None, None, None, False
5995 tool = tools.get(name)
5996 if not tool:
5997 return params, f'Error: Tool "{name}" not found.', None, None, False
5998 spec = tool.get('spec', {})
5999 tool_type = tool.get('type', '')
6000 direct_tool = tool.get('direct', False)
6001 allowed_params = spec.get('parameters', {}).get('properties', {}).keys()
6002 params = {key: value for key, value in params.items() if key in allowed_params}
6003 try:
6004 if direct_tool:
6005 result = await event_caller(
6006 {
6007 'type': 'execute:tool',
6008 'data': {
6009 'id': str(uuid4()),
6010 'name': name,
6011 'params': params,
6012 'server': tool.get('server', {}),
6013 'session_id': metadata.get('session_id'),
6014 },
6015 }
6016 )
6017 else:
6018 function = await get_updated_tool_function(
6019 function=tool['callable'],
6020 extra_params={
6021 '__messages__': form_data.get('messages', []),
6022 '__files__': metadata.get('files', []),
6023 },
6024 )
6025 result = await function(**params)
6026 except Exception as e:
6027 result = {'error': str(e)}
6028 return params, result, tool, tool_type, direct_tool
6030 delegate_calls = [
6031 tool_call
6032 for tool_call in response_tool_calls
6033 if tool_call.get('function', {}).get('name') == 'delegate_task'
6034 ]
6035 tool_results = {}
6036 for tool_call in response_tool_calls:
6037 if tool_call.get('function', {}).get('name') != 'delegate_task':
6038 tool_results[id(tool_call)] = await execute_tool_call(tool_call)
6039 tool_results.update(
6040 zip(
6041 [id(tool_call) for tool_call in delegate_calls],
6042 await asyncio.gather(*(execute_tool_call(tool_call) for tool_call in delegate_calls)),
6043 )
6044 )
6046 for tool_call in response_tool_calls:
6047 tool_call_id = tool_call.get('id', '')
6048 tool_function_name = tool_call.get('function', {}).get('name', '')
6049 tool_function_params, tool_result, tool, tool_type, direct_tool = tool_results[id(tool_call)]
6050 if tool_result is None:
6051 results.append(
6052 {
6053 'tool_call_id': tool_call_id,
6054 'content': (
6055 'Error: Tool call arguments could not be parsed. The model generated '
6056 f'malformed or incomplete JSON for `{tool_function_name}`. Please try again.'
6057 ),
6058 }
6059 )
6060 continue
6062 terminal_file_result = build_terminal_file_tool_result(
6063 tool_function_name,
6064 tool_function_params,
6065 tool_result,
6066 tool,
6067 metadata,
6068 )
6069 if terminal_file_result:
6070 tool_result = terminal_file_result
6072 tool_result, tool_result_files, tool_result_embeds = await process_tool_result(
6073 request,
6074 tool_function_name,
6075 tool_result,
6076 tool_type,
6077 direct_tool,
6078 metadata,
6079 user,
6080 )
6082 await terminal_event_handler(
6083 tool_function_name,
6084 tool_function_params,
6085 tool_result,
6086 event_emitter,
6087 )
6089 # Extract citation sources from tool results
6090 if (
6091 citations_enabled
6092 and tool_function_name
6093 in [
6094 'fetch_url',
6095 'view_file',
6096 'view_knowledge_file',
6097 'query_knowledge_files',
6098 'query_chat_files',
6099 ]
6100 and tool_result
6101 ):
6102 try:
6103 citation_sources = get_citation_source_from_tool_result(
6104 tool_name=tool_function_name,
6105 tool_params=tool_function_params,
6106 tool_result=tool_result,
6107 tool_id=tool.get('tool_id', '') if tool else '',
6108 )
6109 tool_call_sources.extend(citation_sources)
6110 except Exception as e:
6111 log.exception(f'Error extracting citation source: {e}')
6113 results.append(
6114 {
6115 'tool_call_id': tool_call_id,
6116 'content': tool_result_content(tool_result),
6117 **({'files': tool_result_files} if tool_result_files else {}),
6118 **({'embeds': tool_result_embeds} if tool_result_embeds else {}),
6119 }
6120 )
6122 result_status_by_call_id = {}
6123 for result in results:
6124 output_parts = [{'type': 'input_text', 'text': result.get('content', '')}]
6125 local_output_status = (
6126 'failed' if _is_tool_result_error(result.get('content', '')) else 'completed'
6127 )
6128 result_status_by_call_id[result.get('tool_call_id', '')] = local_output_status
6130 # Separate image data URIs (for LLM via input_image) from
6131 # other files (for frontend display via files attribute).
6132 display_files = []
6133 for file_item in result.get('files', []):
6134 if file_item.get('type') == 'image' and file_item.get('url', '').startswith('data:'):
6135 # LLM-only: add as input_image part, not frontend display output.
6136 image_url = await store_tool_result_image(request, file_item['url'], metadata, user)
6137 output_parts.append({'type': 'input_image', 'image_url': image_url})
6138 else:
6139 # Frontend display (MCP images, audio, etc.)
6140 display_files.append(file_item)
6142 output.append(
6143 {
6144 'type': 'function_call_output',
6145 'id': output_id('fco'),
6146 'call_id': result.get('tool_call_id', ''),
6147 'output': output_parts,
6148 'status': local_output_status,
6149 **({'files': display_files} if display_files else {}),
6150 **({'embeds': result.get('embeds')} if result.get('embeds') else {}),
6151 }
6152 )
6154 # Update function_call statuses and parsed/sanitized arguments.
6155 for tc in response_tool_calls:
6156 call_id = tc.get('id', '')
6157 for item in output:
6158 if item.get('type') == 'function_call' and item.get('call_id') == call_id:
6159 item['status'] = result_status_by_call_id.get(call_id, 'completed')
6160 item['arguments'] = tc.get('function', {}).get('arguments', '{}')
6161 break
6163 # Emit citation sources to the frontend for display
6164 if citations_enabled:
6165 for source in tool_call_sources:
6166 await event_emitter({'type': 'source', 'data': source})
6168 # Apply tool source context to messages for the model.
6169 # Restoring to pre-RAG original prevents duplicating
6170 # the RAG template across file and tool sources.
6171 all_tool_call_sources.extend(tool_call_sources)
6172 if all_tool_call_sources and user_message:
6173 # Restore pre-RAG message state before re-applying
6174 # to prevent RAG template duplication.
6175 original_user_message = metadata.get('user_prompt') or user_message
6176 set_last_user_message_content(
6177 original_user_message,
6178 form_data['messages'],
6179 )
6180 if original_system_content is not None:
6181 if get_system_message(form_data['messages']):
6182 replace_system_message_content(
6183 original_system_content,
6184 form_data['messages'],
6185 )
6186 else:
6187 form_data['messages'] = add_or_update_system_message(
6188 original_system_content,
6189 form_data['messages'],
6190 )
6191 else:
6192 replace_system_message_content('', form_data['messages'])
6194 # Build context: file sources with content,
6195 # tool sources as citation markers only.
6196 source_ids = {}
6197 source_context = get_source_context(
6198 metadata.get('sources', []), source_ids
6199 ) + get_source_context(
6200 all_tool_call_sources,
6201 source_ids,
6202 include_content=False,
6203 )
6204 source_context = source_context.strip()
6205 if source_context:
6206 rag_content = await rag_template(
6207 await Config.get('rag.template'),
6208 source_context,
6209 user_message,
6210 )
6211 if RAG_SYSTEM_CONTEXT:
6212 form_data['messages'] = add_or_update_system_message(
6213 rag_content,
6214 form_data['messages'],
6215 append=True,
6216 )
6217 else:
6218 form_data['messages'] = add_or_update_user_message(
6219 rag_content,
6220 form_data['messages'],
6221 append=False,
6222 )
6223 tool_call_sources.clear()
6225 await emit_output()
6227 try:
6228 new_form_data = {
6229 **form_data,
6230 'model': model_id,
6231 'stream': True,
6232 'metadata': metadata,
6233 }
6235 if ENABLE_RESPONSES_API_STATEFUL and last_response_id:
6236 system_message = get_system_message(form_data['messages'])
6237 new_form_data['messages'] = (
6238 [system_message] if system_message else []
6239 ) + convert_output_to_messages(
6240 output, raw=True, reasoning_format=get_reasoning_format(model)
6241 )
6242 new_form_data['previous_response_id'] = last_response_id
6243 else:
6244 tool_messages = convert_output_to_messages(
6245 output,
6246 raw=True,
6247 reasoning_format=get_reasoning_format(model),
6248 flatten_tool_images=True,
6249 )
6251 # Chat Completions providers don't support multimodal
6252 # tool messages. Extract images into a user message.
6253 image_urls = []
6254 for message in tool_messages:
6255 if message.get('role') == 'tool' and isinstance(message.get('content'), list):
6256 text_parts = []
6257 for part in message['content']:
6258 if part.get('type') == 'input_text':
6259 text_parts.append(part.get('text', ''))
6260 elif part.get('type') == 'input_image':
6261 image_urls.append(part.get('image_url', ''))
6262 message['content'] = ''.join(text_parts)
6264 new_form_data['messages'] = [
6265 *form_data['messages'],
6266 *tool_messages,
6267 ]
6269 if image_urls:
6270 new_form_data['messages'].append(
6271 {
6272 'role': 'user',
6273 'content': [
6274 {
6275 'type': 'text',
6276 'text': 'Here are the images from the tool results above. Please analyze them.',
6277 },
6278 *[{'type': 'image_url', 'image_url': {'url': url}} for url in image_urls],
6279 ],
6280 }
6281 )
6283 new_form_data = await convert_url_images_to_base64(new_form_data, user=user)
6285 if filter_functions:
6286 new_form_data, _ = await process_filter_functions(
6287 request=request,
6288 filter_context=filter_context,
6289 filter_functions=filter_functions,
6290 filter_type='request',
6291 form_data=new_form_data,
6292 extra_params=extra_params,
6293 )
6295 new_form_data = normalize_messages_for_model(new_form_data)
6297 res = await generate_chat_completion(
6298 request,
6299 new_form_data,
6300 user,
6301 bypass_system_prompt=True,
6302 )
6304 if isinstance(res, StreamingResponse):
6305 # Save accumulated output and start fresh.
6306 # Responses API output_index values are relative
6307 # to the current response — a clean output list
6308 # keeps indices aligned. The display prefix
6309 # ensures the UI shows tool history during
6310 # streaming.
6311 prior_output = list(full_output())
6312 # Trim the trailing empty placeholder message
6313 # so it doesn't persist as a ghost item once
6314 # the new stream produces real content.
6315 if (
6316 prior_output
6317 and prior_output[-1].get('type') == 'message'
6318 and prior_output[-1].get('status') == 'in_progress'
6319 ):
6320 msg_parts = prior_output[-1].get('content', [])
6321 if not msg_parts or (len(msg_parts) == 1 and not msg_parts[0].get('text', '').strip()):
6322 prior_output.pop()
6323 output = []
6324 output_start = len(prior_output)
6325 await stream_body_handler(res, new_form_data)
6326 output = full_output()
6327 prior_output = []
6328 elif getattr(res, 'status_code', 200) >= 400:
6329 await emit_message_error(get_message_error_content(get_response_error_detail(res)))
6330 break
6331 else:
6332 break
6333 except Exception as e:
6334 error_content = get_message_error_content(e)
6335 log.exception('Tool-call continuation failed: %s', error_content)
6336 await emit_message_error(error_content)
6337 break
6339 if (
6340 max_tool_call_iterations is not None
6341 and tool_calls
6342 and tool_call_iterations >= max_tool_call_iterations
6343 ):
6344 log.warning('Tool-call iteration limit reached (%s)', max_tool_call_iterations)
6345 error_content = f'Tool-call limit reached ({max_tool_call_iterations} iterations).'
6346 await emit_message_error(error_content)
6348 if DETECT_CODE_INTERPRETER:
6349 MAX_RETRIES = 5
6350 retries = 0
6352 while output and output[-1].get('type') == 'open_webui:code_interpreter' and retries < MAX_RETRIES:
6353 await event_emitter(
6354 {
6355 'type': 'chat:completion',
6356 'data': {
6357 'output': full_output(),
6358 },
6359 }
6360 )
6362 retries += 1
6363 log.debug('Attempt count: %s', retries)
6365 ci_item = output[-1]
6366 ci_output = ''
6367 try:
6368 if ci_item.get('attributes', {}).get('type') == 'code':
6369 code = ci_item.get('code', '')
6370 # Sanitize code (strips ANSI codes and markdown fences)
6371 code = sanitize_code(code)
6373 if CODE_INTERPRETER_BLOCKED_MODULES:
6374 blocking_code = textwrap.dedent(f"""
6375 import builtins
6377 BLOCKED_MODULES = {CODE_INTERPRETER_BLOCKED_MODULES}
6379 _real_import = builtins.__import__
6380 def restricted_import(name, globals=None, locals=None, fromlist=(), level=0):
6381 if name.split('.')[0] in BLOCKED_MODULES:
6382 importer_name = globals.get('__name__') if globals else None
6383 if importer_name == '__main__':
6384 raise ImportError(
6385 f"Direct import of module {{name}} is restricted."
6386 )
6387 return _real_import(name, globals, locals, fromlist, level)
6389 builtins.__import__ = restricted_import
6390 """)
6391 code = blocking_code + '\n' + code
6393 ci_engine = await Config.get('code_interpreter.engine')
6394 if ci_engine == 'pyodide':
6395 ci_output = await event_caller(
6396 {
6397 'type': 'execute:python',
6398 'data': {
6399 'id': str(uuid4()),
6400 'code': code,
6401 'session_id': metadata.get('session_id', None),
6402 'files': metadata.get('files', []),
6403 },
6404 }
6405 )
6406 elif ci_engine == 'jupyter':
6407 ci_output = await execute_code_jupyter(
6408 await Config.get('code_interpreter.jupyter.url'),
6409 code,
6410 (
6411 await Config.get('code_interpreter.jupyter.auth_token')
6412 if await Config.get('code_interpreter.jupyter.auth') == 'token'
6413 else None
6414 ),
6415 (
6416 await Config.get('code_interpreter.jupyter.auth_password')
6417 if await Config.get('code_interpreter.jupyter.auth') == 'password'
6418 else None
6419 ),
6420 await Config.get('code_interpreter.jupyter.timeout'),
6421 )
6422 else:
6423 ci_output = {'stdout': 'Code interpreter engine not configured.'}
6425 log.debug('Code interpreter output: %s', ci_output)
6427 # Handle error responses from event_caller
6428 # (e.g. session disconnected, timeout)
6429 if isinstance(ci_output, dict) and ci_output.get('error'):
6430 ci_output = {'stderr': ci_output['error']}
6432 if isinstance(ci_output, dict):
6433 stdout = ci_output.get('stdout', '')
6435 if isinstance(stdout, str):
6436 stdoutLines = stdout.split('\n')
6437 for idx, line in enumerate(stdoutLines):
6438 if re.match(r'data:image/\w+;base64', line):
6439 image_url = await get_image_url_from_base64(
6440 request,
6441 line,
6442 metadata,
6443 user,
6444 )
6445 if image_url:
6446 stdoutLines[idx] = f''
6448 ci_output['stdout'] = '\n'.join(stdoutLines)
6450 result = ci_output.get('result', '')
6452 if isinstance(result, str):
6453 resultLines = result.split('\n')
6454 for idx, line in enumerate(resultLines):
6455 if re.match(r'data:image/\w+;base64', line):
6456 image_url = await get_image_url_from_base64(
6457 request,
6458 line,
6459 metadata,
6460 user,
6461 )
6462 resultLines[idx] = f''
6463 ci_output['result'] = '\n'.join(resultLines)
6464 except Exception as e:
6465 ci_output = str(e)
6467 ci_item['output'] = ci_output
6468 ci_item['status'] = 'completed'
6470 output.append(
6471 {
6472 'type': 'message',
6473 'id': output_id('msg'),
6474 'status': 'in_progress',
6475 'role': 'assistant',
6476 'content': [{'type': 'output_text', 'text': ''}],
6477 }
6478 )
6480 await event_emitter(
6481 {
6482 'type': 'chat:completion',
6483 'data': {
6484 'output': full_output(),
6485 },
6486 }
6487 )
6489 try:
6490 new_form_data = {
6491 **form_data,
6492 'model': model_id,
6493 'stream': True,
6494 'metadata': metadata,
6495 'messages': [
6496 *form_data['messages'],
6497 *convert_output_to_messages(
6498 output,
6499 raw=True,
6500 reasoning_format=get_reasoning_format(model),
6501 flatten_tool_images=True,
6502 ),
6503 ],
6504 }
6506 if filter_functions:
6507 new_form_data, _ = await process_filter_functions(
6508 request=request,
6509 filter_context=filter_context,
6510 filter_functions=filter_functions,
6511 filter_type='request',
6512 form_data=new_form_data,
6513 extra_params=extra_params,
6514 )
6516 new_form_data = normalize_messages_for_model(new_form_data)
6518 res = await generate_chat_completion(
6519 request,
6520 new_form_data,
6521 user,
6522 bypass_system_prompt=True,
6523 )
6525 if isinstance(res, StreamingResponse):
6526 await stream_body_handler(res, new_form_data)
6527 elif getattr(res, 'status_code', 200) >= 400:
6528 await emit_message_error(get_message_error_content(get_response_error_detail(res)))
6529 break
6530 else:
6531 break
6532 except Exception as e:
6533 error_content = get_message_error_content(e)
6534 log.exception('Code interpreter continuation failed: %s', error_content)
6535 await emit_message_error(error_content)
6536 break
6538 # Mark all in-progress items as completed
6539 for item in output:
6540 if item.get('status') == 'in_progress':
6541 item['status'] = 'completed'
6543 current_output = full_output()
6544 title = await Chats.get_chat_title_by_id(metadata['chat_id']) if save_to_chat else ''
6545 data = {
6546 'done': True,
6547 'output': current_output,
6548 'title': title,
6549 **({'usage': usage} if usage else {}),
6550 }
6552 if save_to_chat:
6553 # Save final output once. The delta path keeps in-progress
6554 # state in response_streams instead of writing tokens to DB.
6555 await Chats.upsert_message_to_chat_by_id_and_message_id(
6556 metadata['chat_id'],
6557 metadata['message_id'],
6558 {
6559 'done': True,
6560 'output': current_output,
6561 **({'usage': usage} if usage else {}),
6562 },
6563 )
6565 await clear_response_stream(request.app.state.redis, response_stream_task_id)
6566 await publish_chat_finished_event(
6567 request,
6568 user,
6569 metadata,
6570 title,
6571 get_output_text(current_output) if continuing else ''.join(content_parts),
6572 current_output,
6573 )
6575 await event_emitter(
6576 {
6577 'type': 'chat:completion',
6578 'data': data,
6579 }
6580 )
6582 ctx['assistant_message'] = {
6583 'content': get_output_text(current_output)
6584 if continuing
6585 else ''.join(content_parts) or get_output_text(current_output),
6586 'output': current_output,
6587 **({'usage': usage} if usage else {}),
6588 }
6589 await outlet_filter_handler(ctx)
6590 await background_tasks_handler(ctx)
6591 except asyncio.CancelledError:
6592 log.warning('Task was cancelled!')
6594 # Close the response body iterator to trigger cleanup
6595 # in stream_wrapper's finally block and release the
6596 # upstream connection. Without this, the async
6597 # generator is orphaned and may spin in anyio internals.
6598 if hasattr(response, 'body_iterator') and hasattr(response.body_iterator, 'aclose'):
6599 try:
6600 await asyncio.shield(response.body_iterator.aclose())
6601 except (asyncio.CancelledError, Exception):
6602 pass
6604 async def save_cancelled_state():
6605 cancelled_output = full_output()
6606 result_call_ids = {
6607 item.get('call_id')
6608 for item in cancelled_output
6609 if item.get('type') == 'function_call_output' and item.get('call_id')
6610 }
6611 for item in cancelled_output:
6612 # A tool call is stamped completed when its arguments finish, before the tool runs
6613 is_running_tool_call = (
6614 item.get('type') == 'function_call'
6615 and item.get('status') == 'completed'
6616 and (item.get('call_id') or item.get('id')) not in result_call_ids
6617 )
6618 if item.get('status') == 'in_progress' or is_running_tool_call:
6619 item['status'] = 'incomplete'
6621 await event_emitter({'type': 'chat:tasks:cancel', 'data': {'output': cancelled_output}})
6622 if save_to_chat:
6623 await Chats.upsert_message_to_chat_by_id_and_message_id(
6624 metadata['chat_id'],
6625 metadata['message_id'],
6626 {
6627 'done': True,
6628 'output': cancelled_output,
6629 },
6630 )
6631 await clear_response_stream(request.app.state.redis, response_stream_task_id)
6633 try:
6634 await asyncio.shield(save_cancelled_state())
6635 except (asyncio.CancelledError, Exception):
6636 pass
6637 raise # re-raise CancelledError for proper propagation
6639 if response.background is not None:
6640 await response.background()
6642 return await response_handler(response, events)
6644 else:
6645 # Fallback to the original response
6646 async def stream_wrapper(original_generator, events):
6647 def wrap_item(item):
6648 return f'data: {item}\n\n'
6650 try:
6651 assistant_message = {}
6652 filter_context = FilterContext()
6653 has_api_outlet_filters = ENABLE_API_OUTLET_FILTERS and bool(filter_functions)
6654 if ENABLE_API_OUTLET_FILTERS and not has_api_outlet_filters:
6655 try:
6656 model_id = model.get('id') if isinstance(model, dict) else model
6657 has_api_outlet_filters = bool(
6658 (isinstance(model, dict) and 'pipeline' in model)
6659 or get_sorted_filters(model_id, request.app.state.MODELS)
6660 )
6661 except Exception:
6662 has_api_outlet_filters = True
6664 for event in events:
6665 event, _ = await process_filter_functions(
6666 request=request,
6667 filter_context=filter_context,
6668 filter_functions=filter_functions,
6669 filter_type='stream',
6670 form_data=event,
6671 extra_params=extra_params,
6672 )
6674 if event:
6675 yield wrap_item(JSONCodec.dumps(event))
6677 async for data in original_generator:
6678 if filter_functions:
6679 line = data.decode('utf-8', 'replace') if isinstance(data, bytes) else data
6680 if isinstance(line, str) and line.startswith('data:'):
6681 payload = line.removeprefix('data:').strip()
6682 if payload and payload != '[DONE]':
6683 try:
6684 event = JSONCodec.loads(payload)
6685 except JSONCodec.JSONDecodeError:
6686 event = None
6688 if isinstance(event, dict):
6689 event, _ = await process_filter_functions(
6690 request=request,
6691 filter_context=filter_context,
6692 filter_functions=filter_functions,
6693 filter_type='stream',
6694 form_data=event,
6695 extra_params=extra_params,
6696 )
6697 data = wrap_item(JSONCodec.dumps(event)) if event else None
6699 if data:
6700 if has_api_outlet_filters:
6701 update_assistant_message_from_stream(assistant_message, data)
6702 yield data
6704 if has_api_outlet_filters and assistant_message:
6705 ctx['assistant_message'] = assistant_message
6706 await outlet_filter_handler(ctx)
6707 except Exception as e:
6708 log.exception('Chat completion stream failed mid-response: %s', e)
6709 # Separate the error frame from any unfinished upstream event.
6710 yield f'\n\ndata: {JSONCodec.dumps({"error": {"message": "Chat completion stream failed"}})}\n\n'
6711 yield 'data: [DONE]\n\n'
6713 return StreamingResponse(
6714 stream_wrapper(response.body_iterator, events),
6715 headers=dict(response.headers),
6716 background=response.background,
6717 )
6720async def process_chat_response(response, ctx):
6721 # Non-streaming response
6722 if not isinstance(response, StreamingResponse):
6723 return await non_streaming_chat_response_handler(response, ctx)
6725 # Non standard response
6726 if not any(
6727 content_type in response.headers['Content-Type']
6728 for content_type in ['text/event-stream', 'application/x-ndjson']
6729 ):
6730 return response
6732 # Streaming response
6733 return await streaming_chat_response_handler(response, ctx)