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

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 

20 

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 

153 

154logging.basicConfig(stream=sys.stdout, level=GLOBAL_LOG_LEVEL) 

155log = logging.getLogger(__name__) 

156 

157 

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 

168 

169 parsed = value 

170 while isinstance(parsed, str): 

171 try: 

172 parsed = JSONCodec.loads(parsed) 

173 except (JSONCodec.JSONDecodeError, TypeError, ValueError): 

174 break 

175 

176 if not isinstance(parsed, dict): 

177 return False 

178 

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 

186 

187 status = parsed.get('status') 

188 if isinstance(status, str) and status.strip().lower() in {'error', 'failed'}: 

189 return True 

190 

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 ) 

196 

197 return False 

198 

199 

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 

204 

205 

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 

212 

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}}) 

236 

237 

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] 

252 

253DEFAULT_SOLUTION_TAGS = [('<|begin_of_solution|>', '<|end_of_solution|>')] 

254DEFAULT_CODE_INTERPRETER_TAGS = [('<code_interpreter>', '</code_interpreter>')] 

255 

256 

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) 

261 

262 

263def output_id(prefix: str) -> str: 

264 """Generate OR-style ID: prefix + 24-char hex UUID.""" 

265 return f'{prefix}_{uuid4().hex[:24]}' 

266 

267 

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] 

277 

278 if tool_function_name != 'display_file' or not isinstance(tool_result, dict) or tool_result.get('exists') is False: 

279 return None 

280 

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] 

285 

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') 

294 

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 } 

311 

312 

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) 

319 

320 

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 

327 

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 

334 

335 

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 

341 

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 

351 

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 

357 

358 

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. 

363 

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}' 

367 

368 Each such tool call is split into separate entries so each gets executed 

369 independently. Single-object arguments pass through unchanged. 

370 """ 

371 

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) 

375 

376 decoder = json.JSONDecoder() 

377 results = [] 

378 position = 0 

379 

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] 

391 

392 return results or [raw] 

393 

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) 

402 

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) 

411 

412 return expanded 

413 

414 

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. 

420 

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 

425 

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

435 

436 if tool_name in ('view_knowledge_file', 'view_file'): 

437 if not isinstance(tool_result, dict): 

438 return [] 

439 

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', '') 

444 

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 ] 

463 

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 '') 

468 

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 ] 

482 

483 elif tool_name in ('query_knowledge_files', 'query_chat_files'): 

484 if not isinstance(tool_result, list): 

485 return [] 

486 

487 chunks = tool_result 

488 

489 # Group chunks by source for better citation display 

490 # Each unique source becomes a separate source entry 

491 sources_by_file = {} 

492 

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', '') 

499 

500 # Use file_id or note_id as the key 

501 key = file_id or note_id or source_name 

502 

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 } 

513 

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 ) 

523 

524 # Return all grouped sources as a list 

525 if sources_by_file: 

526 return list(sources_by_file.values()) 

527 

528 # Empty result fallback 

529 return [] 

530 

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 ] 

553 

554 

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 

574 

575 

576RESPONSE_COMPLETION_RESPONSE_FIELDS = ('error', 'id', 'output', 'usage') 

577 

578 

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 

584 

585 response_data = {key: response[key] for key in RESPONSE_COMPLETION_RESPONSE_FIELDS if key in response} 

586 

587 return { 

588 **event, 

589 'response': response_data, 

590 } 

591 

592 

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. 

599 

600 Args: 

601 data: The event data 

602 current_output: List of output items (treated as immutable) 

603 

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. 

612 

613 event_type = data.get('type', '') 

614 

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 

641 

642 elif event_type == 'response.content_part.added': 

643 part = data.get('part', {}) 

644 output_index = data.get('output_index', len(current_output) - 1) 

645 

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 

651 

652 if 'content' not in item: 

653 item['content'] = [] 

654 else: 

655 # Copy content list 

656 item['content'] = list(item['content']) 

657 

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 

665 

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) 

669 

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 

674 

675 if 'summary' not in item: 

676 item['summary'] = [] 

677 else: 

678 item['summary'] = list(item['summary']) 

679 

680 item['summary'].append(part) 

681 return new_output, None 

682 return current_output, None 

683 

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', '') 

690 

691 output_index = data.get('output_index', len(current_output) - 1) 

692 

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', '') 

698 

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 

708 

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 

719 

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'] 

726 

727 while len(content_list) <= content_index: 

728 content_list.append({'type': 'text', 'text': ''}) 

729 

730 # Copy the part to mutate it 

731 part = content_list[content_index].copy() 

732 content_list[content_index] = part 

733 

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 '' 

738 

739 part[key] = deep_merge(current_val, delta) 

740 

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'] 

753 

754 while len(summary_list) <= summary_index: 

755 summary_list.append({'type': 'summary_text', 'text': ''}) 

756 

757 part = summary_list[summary_index].copy() 

758 summary_list[summary_index] = part 

759 

760 target_val = part.get(key, '') 

761 part[key] = deep_merge(target_val, delta) 

762 

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'] 

772 

773 while len(content_list) <= content_index: 

774 # Reasoning content parts default to text 

775 content_list.append({'type': 'text', 'text': ''}) 

776 

777 part = content_list[content_index].copy() 

778 content_list[content_index] = part 

779 

780 target_val = part.get(key, '') 

781 part[key] = deep_merge(target_val, delta) 

782 

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 

788 

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 

795 

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) 

800 

801 return new_output, None 

802 

803 return current_output, None 

804 

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) 

809 

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, {} 

816 

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] 

822 

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) 

830 

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 

835 

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 

843 

844 elif type_name == 'reasoning_summary_part': 

845 part = data.get('part') 

846 output_index = data.get('output_index', len(current_output) - 1) 

847 

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 

852 

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 

860 

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' 

878 

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', '') 

885 

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 

901 

902 return new_output, {} 

903 

904 return current_output, None 

905 

906 elif event_type == 'response.completed': 

907 # State Machine Event: Completed 

908 response_data = data.get('response', {}) 

909 final_output = response_data.get('output') 

910 

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 

913 

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' 

919 

920 return new_output, { 

921 'usage': response_data.get('usage'), 

922 'done': True, 

923 'response_id': response_data.get('id'), 

924 } 

925 

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 

930 

931 elif event_type == 'response.failed': 

932 # State Machine Event: Failed 

933 error = data.get('response', {}).get('error', {}) 

934 return current_output, {'error': error} 

935 

936 else: 

937 return current_output, None 

938 

939 

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 

970 

971 

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. 

982 

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 

989 

990 context = get_source_context(sources, include_content=include_content) 

991 

992 context = context.strip() 

993 if not context: 

994 return messages 

995 

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 ) 

1008 

1009 

1010BASE64_IMAGE_DATA_URI_RE = re.compile(r'data:image/[a-zA-Z0-9.+-]+;base64,[A-Za-z0-9+/]+={0,2}', re.IGNORECASE) 

1011 

1012 

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 

1027 

1028 

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 

1038 

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 

1050 

1051 

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') 

1063 

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 

1070 

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) 

1076 

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') 

1106 

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 

1111 

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) 

1118 

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 ) 

1124 

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 ) 

1134 

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 

1141 

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 

1157 

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 } 

1167 

1168 tool_result_files = [] 

1169 

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.' 

1176 

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 ) 

1202 

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) 

1247 

1248 tool_result = extract_base64_images(tool_result, tool_result_files) 

1249 

1250 if isinstance(tool_result, list): 

1251 tool_result = {'results': tool_result} 

1252 

1253 if isinstance(tool_result, dict) or isinstance(tool_result, list): 

1254 tool_result = json.dumps(tool_result, indent=2, ensure_ascii=False) 

1255 

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) 

1265 

1266 return tool_result, tool_result_files, tool_result_embeds 

1267 

1268 

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 {} 

1275 

1276 

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. 

1284 

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 

1291 

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') 

1304 

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 ) 

1332 

1333 

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'] 

1343 

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 

1350 

1351 def get_tools_function_calling_payload(messages, task_model_id, content): 

1352 user_message = get_last_user_message(messages) 

1353 

1354 if user_message and messages and messages[-1]['role'] == 'user': 

1355 # Remove the last user message to avoid duplication 

1356 messages = messages[:-1] 

1357 

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 ) 

1362 

1363 prompt = f'History:\n{chat_history}\nQuery: {user_message}' if chat_history else f'Query: {user_message}' 

1364 

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 } 

1374 

1375 event_caller = extra_params['__event_call__'] 

1376 event_emitter = extra_params['__event_emitter__'] 

1377 metadata = extra_params['__metadata__'] 

1378 

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 ) 

1391 

1392 skip_files = False 

1393 sources = [] 

1394 

1395 specs = [tool['spec'] for tool in tools.values()] 

1396 tools_specs = JSONCodec.dumps(specs, ensure_ascii=False) 

1397 

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 

1403 

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) 

1406 

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) 

1412 

1413 if not content: 

1414 return body, {} 

1415 

1416 try: 

1417 content = content[content.find('{') : content.rfind('}') + 1] 

1418 if not content: 

1419 raise Exception('No JSON object found in the response') 

1420 

1421 result = JSONCodec.loads(content) 

1422 

1423 async def tool_call_handler(tool_call): 

1424 nonlocal skip_files 

1425 

1426 log.debug('tool_call=%r', tool_call) 

1427 

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 

1432 

1433 tool_function_params = tool_call.get('parameters', {}) 

1434 

1435 tool = None 

1436 tool_type = '' 

1437 direct_tool = False 

1438 

1439 try: 

1440 tool = tools[tool_function_name] 

1441 tool_type = tool.get('type', '') 

1442 direct_tool = tool.get('direct', False) 

1443 

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} 

1447 

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) 

1464 

1465 except Exception as e: 

1466 tool_result = {'error': str(e)} 

1467 

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 ) 

1477 

1478 if event_emitter: 

1479 await terminal_event_handler( 

1480 tool_function_name, 

1481 tool_function_params, 

1482 tool_result, 

1483 event_emitter, 

1484 ) 

1485 

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 ) 

1492 

1493 await event_emitter( 

1494 { 

1495 'type': 'files', 

1496 'data': { 

1497 'files': tool_result_files, 

1498 }, 

1499 } 

1500 ) 

1501 

1502 if tool_result_embeds: 

1503 await event_emitter( 

1504 { 

1505 'type': 'embeds', 

1506 'data': { 

1507 'embeds': tool_result_embeds, 

1508 }, 

1509 } 

1510 ) 

1511 

1512 if tool_result: 

1513 tool = tools[tool_function_name] 

1514 tool_id = tool.get('tool_id', '') 

1515 

1516 tool_name = f'{tool_id}/{tool_function_name}' if tool_id else f'{tool_function_name}' 

1517 

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 ) 

1534 

1535 if tools[tool_function_name].get('metadata', {}).get('file_handler', False): 

1536 skip_files = True 

1537 

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) 

1544 

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 

1551 

1552 log.debug('tool_contexts: %s', sources) 

1553 

1554 if skip_files and 'files' in body.get('metadata', {}): 

1555 del body['metadata']['files'] 

1556 

1557 return body, {'sources': sources} 

1558 

1559 

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 ) 

1572 

1573 messages = form_data['messages'] 

1574 user_message = get_last_user_message(messages) 

1575 

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 ) 

1589 

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) 

1601 

1602 response = res['choices'][0]['message']['content'] 

1603 

1604 try: 

1605 bracket_start = response.rfind('{') 

1606 bracket_end = response.rfind('}') + 1 

1607 

1608 if bracket_start == -1 or bracket_end == -1: 

1609 raise Exception('No JSON object found in the response') 

1610 

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] 

1616 

1617 if ENABLE_QUERIES_CACHE: 

1618 request.state.cached_queries = queries 

1619 

1620 except Exception as e: 

1621 log.exception(e) 

1622 queries = [user_message or ''] 

1623 

1624 # Check if generated queries are empty 

1625 if len(queries) == 1 and queries[0].strip() == '': 

1626 queries = [user_message or ''] 

1627 

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 

1641 

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 ) 

1652 

1653 try: 

1654 results = await process_web_search( 

1655 request, 

1656 SearchForm(queries=queries), 

1657 user=user, 

1658 ) 

1659 

1660 if results: 

1661 files = form_data.get('files', []) 

1662 

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 ) 

1686 

1687 form_data['files'] = files 

1688 

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 ) 

1713 

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 ) 

1729 

1730 return form_data 

1731 

1732 

1733def get_images_from_messages(message_list): 

1734 images = [] 

1735 

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')) 

1743 

1744 if message_images: 

1745 images.append(message_images) 

1746 

1747 return images 

1748 

1749 

1750async def get_image_urls(delta_images, request, metadata, user) -> list[str]: 

1751 if not isinstance(delta_images, list): 

1752 return [] 

1753 

1754 image_urls = [] 

1755 for img in delta_images: 

1756 if not isinstance(img, dict) or img.get('type') != 'image_url': 

1757 continue 

1758 

1759 url = img.get('image_url', {}).get('url') 

1760 if not url: 

1761 continue 

1762 

1763 if url.startswith('data:image/png;base64'): 

1764 url = await get_image_url_from_base64(request, url, metadata, user) 

1765 

1766 image_urls.append(url) 

1767 

1768 return image_urls 

1769 

1770 

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 

1777 

1778 chat = await Chats.get_chat_by_id_and_user_id(chat_id, user.id) 

1779 if not chat: 

1780 return messages 

1781 

1782 history = chat.chat.get('history', {}) 

1783 stored_messages = get_message_list(history.get('messages', {}), history.get('currentId')) 

1784 

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}/>' 

1795 

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'] 

1804 

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 

1815 

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' 

1818 

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 

1824 

1825 return messages 

1826 

1827 

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) 

1832 

1833 if not chat_id or not isinstance(chat_id, str) or not __event_emitter__: 

1834 return form_data 

1835 

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 } 

1841 

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) 

1846 

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) 

1850 

1851 user_message = get_last_user_message(message_list) 

1852 

1853 prompt = user_message 

1854 message_images = get_images_from_messages(message_list) 

1855 

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) 

1864 

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 

1869 

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 ) 

1877 

1878 system_message_content = '' 

1879 

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 ) 

1889 

1890 await __event_emitter__( 

1891 { 

1892 'type': 'status', 

1893 'data': {'description': 'Image created', 'done': True}, 

1894 } 

1895 ) 

1896 

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 ) 

1911 

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) 

1915 

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) 

1922 

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 ) 

1932 

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>' 

1934 

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 ) 

1945 

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>' 

1947 

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 ) 

1961 

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) 

1970 

1971 response = res['choices'][0]['message']['content'] 

1972 

1973 try: 

1974 bracket_start = response.rfind('{') 

1975 bracket_end = response.rfind('}') + 1 

1976 

1977 if bracket_start == -1 or bracket_end == -1: 

1978 raise Exception('No JSON object found in the response') 

1979 

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 

1985 

1986 except Exception as e: 

1987 log.exception(e) 

1988 prompt = user_message 

1989 

1990 try: 

1991 images = await image_generations( 

1992 request=request, 

1993 form_data=CreateImageForm(**{'prompt': prompt}), 

1994 metadata=image_metadata, 

1995 user=user, 

1996 ) 

1997 

1998 await __event_emitter__( 

1999 { 

2000 'type': 'status', 

2001 'data': {'description': 'Image created', 'done': True}, 

2002 } 

2003 ) 

2004 

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 ) 

2019 

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) 

2023 

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) 

2030 

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 ) 

2040 

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>' 

2042 

2043 if system_message_content: 

2044 form_data['messages'] = add_or_update_system_message(system_message_content, form_data['messages']) 

2045 

2046 return form_data 

2047 

2048 

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 = [] 

2054 

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) 

2059 

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'] 

2074 

2075 try: 

2076 bracket_start = queries_response.rfind('{') 

2077 bracket_end = queries_response.rfind('}') + 1 

2078 

2079 if bracket_start == -1 or bracket_end == -1: 

2080 raise Exception('No JSON object found in the response') 

2081 

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]} 

2086 

2087 queries = queries_response.get('queries', []) 

2088 except Exception: 

2089 pass 

2090 

2091 await __event_emitter__( 

2092 { 

2093 'type': 'status', 

2094 'data': { 

2095 'action': 'queries_generated', 

2096 'queries': queries, 

2097 'done': False, 

2098 }, 

2099 } 

2100 ) 

2101 

2102 if len(queries) == 0: 

2103 queries = [get_last_user_message(body['messages']) or ''] 

2104 

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) 

2138 

2139 log.debug('rag_contexts:sources: %s', sources) 

2140 

2141 unique_ids = set() 

2142 for source in sources or []: 

2143 if not source or len(source.keys()) == 0: 

2144 continue 

2145 

2146 documents = source.get('document') or [] 

2147 metadatas = source.get('metadata') or [] 

2148 src_info = source.get('source') or {} 

2149 

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) 

2154 

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 ) 

2166 

2167 return body, {'sources': sources} 

2168 

2169 

2170async def convert_url_images_to_base64(form_data, user=None): 

2171 messages = form_data.get('messages', []) 

2172 

2173 for message in messages: 

2174 content = message.get('content') 

2175 if not isinstance(content, list): 

2176 continue 

2177 

2178 new_content = [] 

2179 

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 

2184 

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 

2195 

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) 

2215 

2216 message['content'] = new_content 

2217 

2218 return form_data 

2219 

2220 

2221MESSAGE_REPLAY_KEYS = ('id', 'role', 'content', 'output', 'files', 'contextSummary', 'usage', 'model') 

2222 

2223 

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 

2232 

2233 db_messages = get_message_list(messages_map, message_id) 

2234 if not db_messages: 

2235 return None 

2236 

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 ] 

2244 

2245 

2246def get_reasoning_format(model: dict) -> str | None: 

2247 """ 

2248 Determine how reasoning should be included in reconstructed messages. 

2249 

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 

2262 

2263 

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 ] 

2269 

2270 

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. 

2277 

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 = [] 

2282 

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 

2297 

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) 

2302 

2303 return processed 

2304 

2305 

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 } 

2312 

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 } 

2319 

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) 

2336 

2337 return sanitized 

2338 

2339 

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). 

2348 

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 

2356 

2357 if not mcp_server_connection: 

2358 log.error(f'MCP server with id {server_id} not found') 

2359 return None 

2360 

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 

2364 

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 ) 

2373 

2374 client = MCPClient() 

2375 await client.connect( 

2376 url=mcp_server_connection.get('url', ''), 

2377 headers=headers if headers else None, 

2378 ) 

2379 

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(',') 

2383 

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)] 

2387 

2388 return client, tool_specs 

2389 

2390 

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'] = '' 

2395 

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 

2399 

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 ] 

2412 

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) 

2422 

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 

2428 

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') 

2431 

2432 form_data = apply_params_to_form_data(form_data, model) 

2433 log.debug('form_data: %s', form_data) 

2434 

2435 # Guided regeneration: extract before it reaches the LLM provider 

2436 regeneration_prompt = form_data.pop('regeneration_prompt', None) 

2437 

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') 

2442 

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}) 

2453 

2454 system_message = get_system_message(form_data.get('messages', [])) 

2455 form_data['messages'] = [system_message, *db_messages] if system_message else db_messages 

2456 

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) 

2480 

2481 if regeneration_prompt: 

2482 form_data['messages'].append({'role': 'user', 'content': regeneration_prompt}) 

2483 

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 

2492 

2493 system_message = get_system_message(form_data.get('messages', [])) 

2494 system_prompt = get_content_from_message(system_message) if system_message else '' 

2495 

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') 

2514 

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) 

2522 

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']) 

2528 

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 

2537 

2538 form_data = await convert_url_images_to_base64(form_data, user=user) 

2539 

2540 event_emitter = await get_event_emitter(metadata) 

2541 event_caller = await get_event_call(metadata) 

2542 

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 

2562 

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 ) 

2569 

2570 events = [] 

2571 sources = [] 

2572 

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) 

2580 

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) 

2584 

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 

2589 

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) 

2603 

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) 

2607 

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 ) 

2619 

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) 

2641 

2642 files = form_data.get('files', []) 

2643 files.extend(knowledge_files) 

2644 form_data['files'] = files 

2645 

2646 variables = form_data.pop('variables', None) 

2647 payload_tools = form_data.get('tools', None) # snapshot before filters 

2648 

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 

2654 

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', [])) 

2660 

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}') 

2671 

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 

2680 

2681 form_data['messages'] = add_or_update_system_message( 

2682 template, 

2683 form_data['messages'], 

2684 ) 

2685 

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) 

2699 

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) 

2710 

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) 

2721 

2722 if 'code_interpreter' in features and features['code_interpreter']: 

2723 engine = await Config.get('code_interpreter.engine', 'pyodide') 

2724 

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 

2730 

2731 # Append filesystem awareness only for pyodide engine 

2732 if engine != 'jupyter': 

2733 prompt += CODE_INTERPRETER_PYODIDE_PROMPT 

2734 

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 ) 

2752 

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 ) 

2765 

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) 

2769 

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']) 

2783 

2784 is_note_chat = bool(chat and (chat.meta or {}).get('internal') is True and (chat.meta or {}).get('type') == 'note') 

2785 

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] 

2808 

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 ) 

2814 

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 ) 

2825 

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)] 

2829 

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)} 

2835 

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) 

2840 

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 ) 

2855 

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 ) 

2861 

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

2876 

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))] 

2882 

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) 

2899 

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 ) 

2906 

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) 

2910 

2911 prompt = get_last_user_message(form_data['messages']) 

2912 

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) 

2927 

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()) 

2932 

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 

2944 

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) 

2953 

2954 log.debug('tool_ids=%r', tool_ids) 

2955 log.debug('direct_tool_servers=%r', direct_tool_servers) 

2956 

2957 tools_dict = {} 

2958 

2959 mcp_clients = {} 

2960 mcp_tools_dict = {} 

2961 

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:') :] 

2968 

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 

2978 

2979 client, tool_specs = result 

2980 mcp_clients[server_id] = client 

2981 

2982 for tool_spec in tool_specs: 

2983 

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 ) 

2990 

2991 return tool_function 

2992 

2993 tool_function = await make_tool_function(client, tool_spec['name']) 

2994 

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) 

3017 

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 ) 

3030 

3031 if mcp_tools_dict: 

3032 tools_dict = {**tools_dict, **mcp_tools_dict} 

3033 

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 

3064 

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 ) 

3076 

3077 tool_specs = tool_server.pop('specs', []) 

3078 

3079 for tool in tool_specs: 

3080 tools_dict[tool['name']] = { 

3081 'spec': tool, 

3082 'direct': True, 

3083 'server': tool_server, 

3084 } 

3085 

3086 if terminal_id and terminal_capability: 

3087 from open_webui.utils.terminals import add_terminal_agents_md, get_terminal_agents_md 

3088 

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) 

3092 

3093 if mcp_clients: 

3094 metadata['mcp_clients'] = mcp_clients 

3095 

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) 

3103 

3104 if (model.get('info', {}).get('meta', {}).get('builtinTools') or {}).get('knowledge', True): 

3105 from html import escape 

3106 

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}/>') 

3117 

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 ) 

3124 

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 

3139 

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 

3179 

3180 for name in shell_tools: 

3181 if not connected or name not in selected: 

3182 tools_dict.pop(name) 

3183 

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 

3188 

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) 

3205 

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) 

3208 

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) 

3215 

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

3233 

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) 

3237 

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 ] 

3244 

3245 if len(sources) > 0: 

3246 events.append({'sources': sources}) 

3247 

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 ) 

3260 

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}') 

3273 

3274 form_data = normalize_messages_for_model(form_data) 

3275 

3276 return form_data, metadata, events 

3277 

3278 

3279async def get_event_emitter_and_caller(metadata): 

3280 event_emitter = None 

3281 event_caller = None 

3282 

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) 

3288 

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) 

3293 

3294 return event_emitter, event_caller 

3295 

3296 

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 } 

3320 

3321 

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) 

3348 

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.'} 

3352 

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} 

3358 

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)} 

3387 

3388 terminal_file_result = build_terminal_file_tool_result(name, params, result, tool, metadata) 

3389 if terminal_file_result: 

3390 result = terminal_file_result 

3391 

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 ) 

3401 

3402 await terminal_event_handler(name, params, result, event_emitter) 

3403 

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 } 

3410 

3411 

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 

3418 

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 

3424 

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 

3453 

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 

3462 

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) 

3490 

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 

3503 

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 ) 

3540 

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 ) 

3557 

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) 

3573 

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']) 

3579 

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) 

3604 

3605 if not paused: 

3606 normalize_messages_for_model(form_data) 

3607 

3608 return paused 

3609 

3610 return False 

3611 

3612 

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'] 

3621 

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' 

3633 

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 ) 

3656 

3657 

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] 

3662 

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 

3675 

3676 return response, response_data 

3677 

3678 

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 

3687 

3688 return { 

3689 **extra_response, 

3690 **response_data, 

3691 } 

3692 return response_data 

3693 

3694 

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 

3705 

3706 

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 

3711 

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}) 

3718 

3719 for raw_part in line.splitlines(): 

3720 part = raw_part.removeprefix('data:').strip() 

3721 if not part or part == '[DONE]': 

3722 continue 

3723 

3724 try: 

3725 data = JSONCodec.loads(part) 

3726 except Exception: 

3727 continue 

3728 

3729 if not isinstance(data, dict): 

3730 continue 

3731 

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 

3739 

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) 

3744 

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') 

3749 

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 ) 

3766 

3767 append_output_text(output[-1], reasoning_content) 

3768 

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']) 

3776 

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 ) 

3787 

3788 append_output_text(output[-1], content) 

3789 

3790 if 'content' in assistant_message: 

3791 append_to_text_field(assistant_message, 'content', content) 

3792 else: 

3793 assistant_message['content'] = '' + content 

3794 

3795 

3796async def get_system_oauth_token(request, user): 

3797 """Get the system OAuth token for a user. 

3798 

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 ) 

3811 

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 

3815 

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 

3830 

3831 

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'] 

3839 

3840 message = None 

3841 messages = [] 

3842 

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']) 

3849 

3850 message_list = get_message_list(messages_map, metadata['message_id']) 

3851 

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 

3855 

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 

3864 

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() 

3872 

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') 

3886 

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 ) 

3900 

3901 if res and isinstance(res, dict): 

3902 if len(res.get('choices', [])) == 1: 

3903 response_message = res.get('choices', [])[0].get('message', {}) 

3904 

3905 follow_ups_string = response_message.get('content') or response_message.get( 

3906 'reasoning_content', '' 

3907 ) 

3908 else: 

3909 follow_ups_string = '' 

3910 

3911 follow_ups_string = follow_ups_string[ 

3912 follow_ups_string.find('{') : follow_ups_string.rfind('}') + 1 

3913 ] 

3914 

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 ) 

3925 

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 ) 

3935 

3936 except Exception as e: 

3937 pass 

3938 

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] + '...' 

3944 

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 ) 

3956 

3957 if res and isinstance(res, dict): 

3958 if len(res.get('choices', [])) == 1: 

3959 response_message = res.get('choices', [])[0].get('message', {}) 

3960 

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 = '' 

3970 

3971 title_string = title_string[title_string.find('{') : title_string.rfind('}') + 1] 

3972 

3973 try: 

3974 title = JSONCodec.loads(title_string).get('title', user_message) 

3975 except Exception as e: 

3976 title = '' 

3977 

3978 if not title: 

3979 title = messages[0].get('content', user_message) 

3980 

3981 await Chats.update_chat_title_by_id(metadata['chat_id'], title) 

3982 

3983 await event_emitter( 

3984 { 

3985 'type': 'chat:title', 

3986 'data': title, 

3987 } 

3988 ) 

3989 

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) 

3992 

3993 await Chats.update_chat_title_by_id(metadata['chat_id'], title) 

3994 

3995 await event_emitter( 

3996 { 

3997 'type': 'chat:title', 

3998 'data': message.get('content', user_message), 

3999 } 

4000 ) 

4001 

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 ) 

4012 

4013 if res and isinstance(res, dict): 

4014 if len(res.get('choices', [])) == 1: 

4015 response_message = res.get('choices', [])[0].get('message', {}) 

4016 

4017 tags_string = response_message.get('content') or response_message.get( 

4018 'reasoning_content', '' 

4019 ) 

4020 else: 

4021 tags_string = '' 

4022 

4023 tags_string = tags_string[tags_string.find('{') : tags_string.rfind('}') + 1] 

4024 

4025 try: 

4026 tags = JSONCodec.loads(tags_string).get('tags', []) 

4027 await Chats.update_chat_tags_by_id(metadata['chat_id'], tags, user) 

4028 

4029 await event_emitter( 

4030 { 

4031 'type': 'chat:tags', 

4032 'data': tags, 

4033 } 

4034 ) 

4035 except Exception as e: 

4036 pass 

4037 

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 ) 

4048 

4049 

4050async def outlet_filter_handler(ctx): 

4051 """Run outlet filters inline after chat completion. 

4052 

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. 

4057 

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') 

4066 

4067 chat_id = metadata.get('chat_id', '') 

4068 message_id = metadata.get('message_id') 

4069 

4070 if not chat_id and not ctx.get('assistant_message'): 

4071 return 

4072 

4073 if not message_id: 

4074 message_id = output_id('msg') 

4075 

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 

4088 

4089 messages_map = None 

4090 

4091 if is_unsaved_chat: 

4092 form_messages = ctx.get('form_data', {}).get('messages', []) 

4093 assistant_message = ctx.get('assistant_message', {}) 

4094 

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 ] 

4102 

4103 if assistant_message: 

4104 message_list.append( 

4105 { 

4106 'id': message_id, 

4107 'role': 'assistant', 

4108 **assistant_message, 

4109 } 

4110 ) 

4111 

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 

4118 

4119 message_list = get_message_list(messages_map, message_id) 

4120 if not message_list: 

4121 return 

4122 

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 } 

4144 

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) 

4150 

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 } 

4160 

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 

4172 

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 ) 

4199 

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) 

4209 

4210 

4211async def non_streaming_chat_response_handler(response, ctx): 

4212 request = ctx['request'] 

4213 

4214 user = ctx['user'] 

4215 metadata = ctx['metadata'] 

4216 events = ctx['events'] 

4217 

4218 event_emitter = ctx['event_emitter'] 

4219 

4220 response, response_data = get_response_data(response) 

4221 if response_data is None: 

4222 return response 

4223 

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')) 

4227 

4228 if event_emitter: 

4229 try: 

4230 if 'error' in response_data: 

4231 error = response_data.get('error') 

4232 

4233 if isinstance(error, dict): 

4234 error = error.get('detail', error) 

4235 else: 

4236 error = str(error) 

4237 

4238 log.error('Provider returned error (non-streaming): %s', error) 

4239 

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 ) 

4256 

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 ) 

4266 

4267 choices = response_data.get('choices', []) 

4268 response_output = response_data.get('output') 

4269 content = choices[0].get('message', {}).get('content') if choices else '' 

4270 

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 ) 

4282 

4283 title = await Chats.get_chat_title_by_id(metadata['chat_id']) if save_to_chat else '' 

4284 

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 ) 

4319 

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) 

4349 

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 ) 

4361 

4362 # Save message in the database 

4363 usage = normalize_usage(response_data.get('usage', {}) or {}) 

4364 

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 ) 

4376 

4377 await publish_chat_finished_event(request, user, metadata, title, content, response_output) 

4378 

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) 

4386 

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 

4410 

4411 return response 

4412 

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) 

4424 

4425 if isinstance(response, dict): 

4426 response = merge_events_into_response(response_data, events) 

4427 

4428 return response 

4429 

4430 

4431async def streaming_chat_response_handler(response, ctx): 

4432 request = ctx['request'] 

4433 

4434 form_data = ctx['form_data'] 

4435 

4436 user = ctx['user'] 

4437 model = ctx['model'] 

4438 

4439 metadata = ctx['metadata'] 

4440 events = ctx['events'] 

4441 

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')) 

4447 

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 } 

4459 

4460 filter_functions = ( 

4461 await get_filter_functions(request, model, metadata.get('filter_ids', [])) if ENABLE_PLUGINS else [] 

4462 ) 

4463 

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', '') 

4470 

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') 

4477 

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. 

4483 

4484 Uses the text from the output items themselves for tag detection, 

4485 eliminating state divergence between accumulated content and items. 

4486 

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 

4491 

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 

4501 

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 '' 

4509 

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 

4516 

4517 def get_scanned_length(item, text): 

4518 item_id = item.get('id') 

4519 if not item_id: 

4520 return 0 

4521 

4522 scanned_length = tag_scan_positions.get((item_id, content_type), 0) 

4523 return scanned_length if scanned_length <= len(text) else 0 

4524 

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) 

4529 

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) 

4535 

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 

4542 

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) 

4555 

4556 return last_open, last_boundary 

4557 

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) 

4565 

4566 last_type = output[-1].get('type', '') if output else '' 

4567 

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) 

4575 

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) 

4582 

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 = '' 

4591 

4592 attributes = extract_attributes(attr_content) 

4593 

4594 before_tag = item_text[: match.start()] 

4595 after_tag = item_text[match.end() :] 

4596 

4597 # Keep only text before the tag in the message 

4598 set_last_text(output, before_tag) 

4599 

4600 if not before_tag.strip(): 

4601 # Remove empty message item 

4602 if output and output[-1].get('type') == 'message': 

4603 output.pop() 

4604 

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 ) 

4651 

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) 

4660 

4661 _, recursive_end = tag_output_handler(content_type, tags, output) 

4662 if recursive_end: 

4663 end_flag = True 

4664 

4665 return output, end_flag 

4666 else: 

4667 save_scanned_length(item, item_text) 

4668 

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', '') 

4677 

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) 

4688 

4689 scanned_length = get_scanned_length(item, block_content) 

4690 end_tag_search_start = max(0, scanned_length - max(len(end_tag), 1) + 1) 

4691 

4692 if block_content.find(end_tag, end_tag_search_start) != -1: 

4693 clear_scanned_length(item) 

4694 end_flag = True 

4695 

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() 

4699 

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) 

4703 

4704 block_content = split_content[0].strip() if split_content else '' 

4705 leftover_content = split_content[1].strip() if len(split_content) > 1 else '' 

4706 

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() 

4721 

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) 

4757 

4758 return None, False 

4759 

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 ) 

4769 

4770 tool_calls = [] 

4771 

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 

4778 

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

4783 

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 = [] 

4815 

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 = [] 

4820 

4821 usage = None 

4822 last_response_id = None 

4823 

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 

4844 

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) 

4852 

4853 return error if isinstance(error, (str, dict)) else str(error) 

4854 

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 ) 

4868 

4869 reasoning_tags_param = metadata.get('params', {}).get('reasoning_tags') 

4870 DETECT_REASONING_TAGS = reasoning_tags_param is not False 

4871 

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 ) 

4892 

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 

4899 

4900 try: 

4901 for event in events: 

4902 await event_emitter( 

4903 { 

4904 'type': 'chat:completion', 

4905 'data': event, 

4906 } 

4907 ) 

4908 

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 ) 

4918 

4919 async def stream_body_handler(response, form_data): 

4920 nonlocal usage 

4921 nonlocal output 

4922 nonlocal prior_output 

4923 nonlocal last_response_id 

4924 

4925 response_tool_calls = [] 

4926 

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 

4935 

4936 joined_content = '' 

4937 joined_part_count = 0 

4938 

4939 async def save_current_response_stream(stream_output: list | None = None): 

4940 nonlocal joined_content 

4941 nonlocal joined_part_count 

4942 

4943 if not chat_id or not metadata.get('message_id'): 

4944 return 

4945 

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) 

4950 

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 ) 

4962 

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 ) 

4974 

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 

4990 

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 

4996 

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 

5011 

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 

5017 

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() 

5031 

5032 delta_count += 1 

5033 last_delta_data = delta_data 

5034 last_delta_type = delta_type 

5035 last_delta_key = delta_key 

5036 

5037 if delta_count >= delta_chunk_size: 

5038 await flush_pending_delta_data(delta_chunk_size) 

5039 

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 

5047 

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) 

5059 

5060 filter_extra_params = {'__body__': form_data, **extra_params} if filter_functions else None 

5061 

5062 async for line in response.body_iterator: 

5063 line = line.decode('utf-8', 'replace') if isinstance(line, bytes) else line 

5064 data = line 

5065 

5066 # Skip empty lines 

5067 if not data or data.isspace(): 

5068 continue 

5069 

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 

5094 

5095 # Remove the "data:" prefix 

5096 data = data[5:].strip() 

5097 

5098 try: 

5099 data = JSONCodec.loads(data) 

5100 

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 ) 

5110 

5111 if data: 

5112 if 'event' in data and not getattr(request.state, 'direct', False): 

5113 await event_emitter(data.get('event', {})) 

5114 

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) 

5137 

5138 if not response_data_is_delta: 

5139 await flush_pending_delta_data() 

5140 

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) 

5151 

5152 url = url_citation.get('url', '') 

5153 title = url_citation.get('title', url) 

5154 

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 ) 

5174 

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 

5185 

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 

5190 

5191 if response_metadata.get('error'): 

5192 await event_emitter( 

5193 { 

5194 'type': 'chat:completion', 

5195 'data': {'error': response_metadata['error']}, 

5196 } 

5197 ) 

5198 

5199 await emit_response_completion_event(data) 

5200 

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', []) 

5211 

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 ) 

5225 

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 

5250 

5251 delta = choices[0].get('delta', {}) 

5252 delta_type = 'content' 

5253 

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'] 

5263 

5264 url = url_citation.get('url', '') 

5265 title = url_citation.get('title', url) 

5266 

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 ) 

5285 

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') 

5290 

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 

5298 

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 ) 

5319 

5320 if delta_name: 

5321 current_response_tool_call['function']['name'] = delta_name 

5322 

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 ) 

5337 

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 } 

5345 

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', '') 

5375 

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 ) 

5402 

5403 await save_current_response_stream() 

5404 data = None 

5405 delta_type = 'tool_call' 

5406 

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 

5424 

5425 await event_emitter( 

5426 { 

5427 'type': 'files', 

5428 'data': {'files': message_files}, 

5429 } 

5430 ) 

5431 

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}' 

5436 

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 ) 

5476 

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] 

5499 

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 ] 

5512 

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' 

5525 

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 

5538 

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' 

5551 

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 ) 

5566 

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 ) 

5577 

5578 # closure-cell str += recopies per chunk; append + join once at read is O(n) 

5579 content_parts.append(value) 

5580 

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 ) 

5603 

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 ) 

5647 

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 ] 

5659 

5660 tag_output = None 

5661 

5662 if DETECT_REASONING_TAGS: 

5663 tag_output, _ = tag_output_handler( 

5664 'reasoning', 

5665 reasoning_tags, 

5666 output, 

5667 ) 

5668 

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 

5676 

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 

5685 

5686 if end: 

5687 break 

5688 

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 

5706 

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 

5715 

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() 

5735 

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() 

5742 

5743 if not parts[-1]['text']: 

5744 output.pop() 

5745 

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 ) 

5756 

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' 

5766 

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)) 

5792 

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)) 

5822 

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() 

5829 

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']) 

5839 

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 ) 

5844 

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 ) 

5854 

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) 

5869 

5870 if snapshot: 

5871 await event_emitter( 

5872 { 

5873 'type': 'chat:completion', 

5874 'data': {'output': frontend_output, 'flush': True}, 

5875 } 

5876 ) 

5877 return 

5878 

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 ) 

5890 

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 

5895 

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 

5915 

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 ) 

5933 

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 

5957 

5958 await emit_output() 

5959 

5960 tools = metadata.get('tools', {}) 

5961 

5962 results = [] 

5963 

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 

5980 

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 

6029 

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 ) 

6045 

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 

6061 

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 

6071 

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 ) 

6081 

6082 await terminal_event_handler( 

6083 tool_function_name, 

6084 tool_function_params, 

6085 tool_result, 

6086 event_emitter, 

6087 ) 

6088 

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}') 

6112 

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 ) 

6121 

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 

6129 

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) 

6141 

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 ) 

6153 

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 

6162 

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}) 

6167 

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']) 

6193 

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() 

6224 

6225 await emit_output() 

6226 

6227 try: 

6228 new_form_data = { 

6229 **form_data, 

6230 'model': model_id, 

6231 'stream': True, 

6232 'metadata': metadata, 

6233 } 

6234 

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 ) 

6250 

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) 

6263 

6264 new_form_data['messages'] = [ 

6265 *form_data['messages'], 

6266 *tool_messages, 

6267 ] 

6268 

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 ) 

6282 

6283 new_form_data = await convert_url_images_to_base64(new_form_data, user=user) 

6284 

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 ) 

6294 

6295 new_form_data = normalize_messages_for_model(new_form_data) 

6296 

6297 res = await generate_chat_completion( 

6298 request, 

6299 new_form_data, 

6300 user, 

6301 bypass_system_prompt=True, 

6302 ) 

6303 

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 

6338 

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) 

6347 

6348 if DETECT_CODE_INTERPRETER: 

6349 MAX_RETRIES = 5 

6350 retries = 0 

6351 

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 ) 

6361 

6362 retries += 1 

6363 log.debug('Attempt count: %s', retries) 

6364 

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) 

6372 

6373 if CODE_INTERPRETER_BLOCKED_MODULES: 

6374 blocking_code = textwrap.dedent(f""" 

6375 import builtins 

6376  

6377 BLOCKED_MODULES = {CODE_INTERPRETER_BLOCKED_MODULES} 

6378  

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) 

6388  

6389 builtins.__import__ = restricted_import 

6390 """) 

6391 code = blocking_code + '\n' + code 

6392 

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.'} 

6424 

6425 log.debug('Code interpreter output: %s', ci_output) 

6426 

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']} 

6431 

6432 if isinstance(ci_output, dict): 

6433 stdout = ci_output.get('stdout', '') 

6434 

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'![Output Image]({image_url})' 

6447 

6448 ci_output['stdout'] = '\n'.join(stdoutLines) 

6449 

6450 result = ci_output.get('result', '') 

6451 

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'![Output Image]({image_url})' 

6463 ci_output['result'] = '\n'.join(resultLines) 

6464 except Exception as e: 

6465 ci_output = str(e) 

6466 

6467 ci_item['output'] = ci_output 

6468 ci_item['status'] = 'completed' 

6469 

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 ) 

6479 

6480 await event_emitter( 

6481 { 

6482 'type': 'chat:completion', 

6483 'data': { 

6484 'output': full_output(), 

6485 }, 

6486 } 

6487 ) 

6488 

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 } 

6505 

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 ) 

6515 

6516 new_form_data = normalize_messages_for_model(new_form_data) 

6517 

6518 res = await generate_chat_completion( 

6519 request, 

6520 new_form_data, 

6521 user, 

6522 bypass_system_prompt=True, 

6523 ) 

6524 

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 

6537 

6538 # Mark all in-progress items as completed 

6539 for item in output: 

6540 if item.get('status') == 'in_progress': 

6541 item['status'] = 'completed' 

6542 

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 } 

6551 

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 ) 

6564 

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 ) 

6574 

6575 await event_emitter( 

6576 { 

6577 'type': 'chat:completion', 

6578 'data': data, 

6579 } 

6580 ) 

6581 

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!') 

6593 

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 

6603 

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' 

6620 

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) 

6632 

6633 try: 

6634 await asyncio.shield(save_cancelled_state()) 

6635 except (asyncio.CancelledError, Exception): 

6636 pass 

6637 raise # re-raise CancelledError for proper propagation 

6638 

6639 if response.background is not None: 

6640 await response.background() 

6641 

6642 return await response_handler(response, events) 

6643 

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' 

6649 

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 

6663 

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 ) 

6673 

6674 if event: 

6675 yield wrap_item(JSONCodec.dumps(event)) 

6676 

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 

6687 

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 

6698 

6699 if data: 

6700 if has_api_outlet_filters: 

6701 update_assistant_message_from_stream(assistant_message, data) 

6702 yield data 

6703 

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' 

6712 

6713 return StreamingResponse( 

6714 stream_wrapper(response.body_iterator, events), 

6715 headers=dict(response.headers), 

6716 background=response.background, 

6717 ) 

6718 

6719 

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) 

6724 

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 

6731 

6732 # Streaming response 

6733 return await streaming_chat_response_handler(response, ctx)