Coverage for open_webui/tools/builtin.py: 4%

1844 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 05:07 +0000

1""" 

2Built-in tools for Open WebUI. 

3 

4These tools are automatically available when native function calling is enabled. 

5 

6IMPORTANT: DO NOT IMPORT THIS MODULE DIRECTLY IN OTHER PARTS OF THE CODEBASE. 

7""" 

8 

9import asyncio 

10import logging 

11import time 

12from typing import Literal, Optional 

13from urllib.parse import unquote 

14 

15from fastapi import HTTPException, Request 

16 

17from open_webui.config import RAG_EMBEDDING_QUERY_PREFIX 

18from open_webui.env import ( 

19 KNOWLEDGE_GREP_MAX_MATCHES, 

20 VIEW_FILE_DEFAULT_MAX_CHARS, 

21 VIEW_FILE_MAX_CHARS, 

22) 

23from open_webui.events import EVENTS, publish_event 

24from open_webui.models.channels import Channel, ChannelMember, Channels 

25from open_webui.models.chats import Chats, chat_search_content_query, chat_search_terms 

26from open_webui.models.config import Config 

27from open_webui.models.groups import Groups 

28from open_webui.models.memories import Memories 

29from open_webui.models.messages import Message, Messages 

30from open_webui.models.notes import Notes 

31from open_webui.models.users import UserModel 

32from open_webui.retrieval.utils import filter_source_metadata, get_content_from_url 

33from open_webui.retrieval.vector.async_client import ASYNC_VECTOR_DB_CLIENT 

34from open_webui.routers.images import ( 

35 CreateImageForm, 

36 EditImageForm, 

37 image_edits, 

38 image_generations, 

39) 

40from open_webui.routers.memories import ( 

41 AddMemoryForm, 

42 ListMemoryPathsForm, 

43 MemoryUpdateModel, 

44 ReadMemoryPathForm, 

45 SearchMemoriesForm, 

46 UpdateMemoriesForm, 

47 update_memory_by_id, 

48) 

49from open_webui.routers.memories import ( 

50 add_memory as _add_memory, 

51) 

52from open_webui.routers.memories import ( 

53 list_memory_paths as _list_memory_paths, 

54) 

55from open_webui.routers.memories import ( 

56 read_memory_path as _read_memory_path, 

57) 

58from open_webui.routers.memories import ( 

59 search_memories as _search_memories, 

60) 

61from open_webui.routers.memories import ( 

62 update_memories as _update_memories, 

63) 

64from open_webui.routers.retrieval import search_web as _search_web 

65from open_webui.socket.main import sio 

66from open_webui.tasks import stop_item_tasks 

67from open_webui.tools.knowledge_fs import kb_exec # noqa: F401 — re-exported 

68from open_webui.utils.chat_id import is_saved_chat_id 

69from open_webui.utils.json_codec import JSONCodec 

70from open_webui.utils.notifications import notify_target 

71from open_webui.utils.sanitize import sanitize_code 

72 

73log = logging.getLogger(__name__) 

74 

75MAX_KNOWLEDGE_BASE_SEARCH_ITEMS = 10_000 

76 

77 

78async def _has_write_access_to_note(note, user_id: str) -> bool: 

79 if note.user_id == user_id: 

80 return True 

81 

82 from open_webui.models.access_grants import AccessGrants 

83 

84 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

85 return await AccessGrants.has_access( 

86 user_id=user_id, 

87 resource_type='note', 

88 resource_id=note.id, 

89 permission='write', 

90 user_group_ids=set(user_group_ids), 

91 ) 

92 

93 

94async def _emit_note_updated(request: Request, user: dict, note) -> None: 

95 await sio.emit('events:note', note.model_dump(), to=f'note:{note.id}') 

96 await publish_event( 

97 request, 

98 EVENTS.NOTE_UPDATED, 

99 actor=user, 

100 subject_id=note.id, 

101 data={'title': note.title}, 

102 ) 

103 

104 

105async def _has_read_access_to_file( 

106 file, 

107 user: dict, 

108 model_knowledge: Optional[list[dict]] = None, 

109) -> bool: 

110 """Check if a user can read a file via ownership, admin role, model attachment, or access grants.""" 

111 user_id = user.get('id') 

112 user_role = user.get('role', 'user') 

113 if file.user_id == user_id or user_role == 'admin': 

114 return True 

115 if model_knowledge and any(item.get('type') == 'file' and item.get('id') == file.id for item in model_knowledge): 

116 return True 

117 from open_webui.utils.access_control.files import has_access_to_file 

118 

119 return await has_access_to_file( 

120 file_id=file.id, 

121 access_type='read', 

122 user=UserModel(**user), 

123 ) 

124 

125 

126# ============================================================================= 

127# TIME UTILITIES 

128# ============================================================================= 

129 

130 

131async def notify( 

132 message: str, 

133 target: str = '', 

134 title: str = '', 

135 __request__: Request = None, 

136 __user__: dict = None, 

137) -> str: 

138 """ 

139 Send a notification to the user's configured notification target. 

140 

141 :param message: Notification body. 

142 :param target: Optional target id or name. Empty uses the default target. 

143 :param title: Optional notification title. 

144 """ 

145 user_id = (__user__ or {}).get('id') 

146 if not user_id: 

147 return 'Notification failed: user not found.' 

148 

149 app_name = getattr(getattr(__request__, 'app', None), 'state', None) 

150 # LICENSE covers this Open WebUI notification identifier. 

151 # Do not alter, remove, obscure, or replace it except as LICENSE permits: 

152 # https://docs.openwebui.com/license. 

153 app_name = getattr(app_name, 'WEBUI_NAME', 'Open WebUI') 

154 try: 

155 result = await notify_target(user_id, message, target=target, title=title, app_name=app_name) 

156 return f'Notification sent to {result.get("target_id")}.' 

157 except Exception as e: 

158 return f'Notification failed: {e}' 

159 

160 

161async def get_current_timestamp( 

162 __request__: Request = None, 

163 __user__: dict = None, 

164) -> str: 

165 """ 

166 Get the current Unix timestamp in seconds. 

167 

168 :return: JSON with current_timestamp (seconds), current_iso (UTC ISO format), and user_local_iso (user's local time) 

169 """ 

170 try: 

171 import datetime 

172 from zoneinfo import ZoneInfo 

173 

174 now = datetime.datetime.now(datetime.timezone.utc) 

175 result = { 

176 'current_timestamp': int(now.timestamp()), 

177 'current_iso': now.isoformat(), 

178 } 

179 

180 # Include the user's local time if timezone is available 

181 tz_name = __user__.get('timezone') if __user__ else None 

182 if tz_name: 

183 try: 

184 user_tz = ZoneInfo(tz_name) 

185 user_now = now.astimezone(user_tz) 

186 result['user_local_iso'] = user_now.isoformat() 

187 result['user_timezone'] = tz_name 

188 except Exception: 

189 pass 

190 

191 return JSONCodec.dumps(result, ensure_ascii=False) 

192 except Exception as e: 

193 log.exception(f'get_current_timestamp error: {e}') 

194 return JSONCodec.dumps({'error': str(e)}) 

195 

196 

197async def calculate_timestamp( 

198 days_ago: int = 0, 

199 weeks_ago: int = 0, 

200 months_ago: int = 0, 

201 years_ago: int = 0, 

202 __request__: Request = None, 

203 __user__: dict = None, 

204) -> str: 

205 """ 

206 Get the current Unix timestamp, optionally adjusted by days, weeks, months, or years. 

207 Use this to calculate timestamps for date filtering in search functions. 

208 Examples: "last week" = weeks_ago=1, "3 days ago" = days_ago=3, "a year ago" = years_ago=1 

209 

210 :param days_ago: Number of days to subtract from current time (default: 0) 

211 :param weeks_ago: Number of weeks to subtract from current time (default: 0) 

212 :param months_ago: Number of months to subtract from current time (default: 0) 

213 :param years_ago: Number of years to subtract from current time (default: 0) 

214 :return: JSON with current_timestamp and calculated_timestamp (both in seconds) 

215 """ 

216 try: 

217 import datetime 

218 

219 from dateutil.relativedelta import relativedelta 

220 

221 now = datetime.datetime.now(datetime.timezone.utc) 

222 current_ts = int(now.timestamp()) 

223 

224 # Calculate the adjusted time 

225 total_days = days_ago + (weeks_ago * 7) 

226 adjusted = now - datetime.timedelta(days=total_days) 

227 

228 # Handle months and years separately (variable length) 

229 if months_ago > 0 or years_ago > 0: 

230 adjusted = adjusted - relativedelta(months=months_ago, years=years_ago) 

231 

232 adjusted_ts = int(adjusted.timestamp()) 

233 

234 result = { 

235 'current_timestamp': current_ts, 

236 'current_iso': now.isoformat(), 

237 'calculated_timestamp': adjusted_ts, 

238 'calculated_iso': adjusted.isoformat(), 

239 } 

240 

241 # Include the user's local time if timezone is available 

242 tz_name = __user__.get('timezone') if __user__ else None 

243 if tz_name: 

244 try: 

245 from zoneinfo import ZoneInfo 

246 

247 user_tz = ZoneInfo(tz_name) 

248 result['user_local_iso'] = now.astimezone(user_tz).isoformat() 

249 result['calculated_local_iso'] = adjusted.astimezone(user_tz).isoformat() 

250 result['user_timezone'] = tz_name 

251 except Exception: 

252 pass 

253 

254 return JSONCodec.dumps(result, ensure_ascii=False) 

255 except ImportError: 

256 # Fallback without dateutil 

257 import datetime 

258 

259 now = datetime.datetime.now(datetime.timezone.utc) 

260 current_ts = int(now.timestamp()) 

261 total_days = days_ago + (weeks_ago * 7) + (months_ago * 30) + (years_ago * 365) 

262 adjusted = now - datetime.timedelta(days=total_days) 

263 adjusted_ts = int(adjusted.timestamp()) 

264 result = { 

265 'current_timestamp': current_ts, 

266 'current_iso': now.isoformat(), 

267 'calculated_timestamp': adjusted_ts, 

268 'calculated_iso': adjusted.isoformat(), 

269 } 

270 

271 tz_name = __user__.get('timezone') if __user__ else None 

272 if tz_name: 

273 try: 

274 from zoneinfo import ZoneInfo 

275 

276 user_tz = ZoneInfo(tz_name) 

277 result['user_local_iso'] = now.astimezone(user_tz).isoformat() 

278 result['calculated_local_iso'] = adjusted.astimezone(user_tz).isoformat() 

279 result['user_timezone'] = tz_name 

280 except Exception: 

281 pass 

282 

283 return JSONCodec.dumps(result, ensure_ascii=False) 

284 except Exception as e: 

285 log.exception(f'calculate_timestamp error: {e}') 

286 return JSONCodec.dumps({'error': str(e)}) 

287 

288 

289# ============================================================================= 

290# WEB SEARCH TOOLS 

291# ============================================================================= 

292 

293 

294async def search_web( 

295 query: str, 

296 count: Optional[int] = None, 

297 __request__: Request = None, 

298 __user__: dict = None, 

299) -> str: 

300 """ 

301 Search the public web for information. Best for current events, external references, 

302 or topics not covered in internal documents. 

303 

304 :param query: The search query to look up 

305 :param count: Number of results to return (default: admin-configured value) 

306 :return: JSON with search results containing title, link, and snippet for each result 

307 """ 

308 if __request__ is None: 

309 return JSONCodec.dumps({'error': 'Request context not available'}) 

310 

311 try: 

312 engine = await Config.get('web.search.engine') 

313 user = UserModel(**__user__) if __user__ else None 

314 

315 configured = await Config.get('web.search.result_count') 

316 max_count = 5 if configured is None else configured 

317 count = max(1, min(count, max_count)) if count is not None else max_count 

318 

319 results = await _search_web(__request__, engine, query, user) 

320 

321 # Limit results 

322 results = results[:count] if results else [] 

323 

324 return JSONCodec.dumps( 

325 [{'title': r.title, 'link': r.link, 'snippet': r.snippet} for r in results], 

326 ensure_ascii=False, 

327 ) 

328 except Exception as e: 

329 log.exception(f'search_web error: {e}') 

330 return JSONCodec.dumps({'error': str(e)}) 

331 

332 

333async def fetch_url( 

334 url: str, 

335 __request__: Request = None, 

336 __user__: dict = None, 

337) -> str: 

338 """ 

339 Fetch and extract the main text content from a web page URL. 

340 

341 :param url: The URL to fetch content from 

342 :return: The extracted text content from the page 

343 """ 

344 if __request__ is None: 

345 return JSONCodec.dumps({'error': 'Request context not available'}) 

346 

347 try: 

348 content, _ = await get_content_from_url(__request__, url) 

349 

350 # Truncate if configured (WEB_FETCH_MAX_CONTENT_LENGTH) 

351 # Guard: content may be None if the web loader silently failed 

352 if content is not None: 

353 max_length = await Config.get('web.fetch.max_content_length') 

354 if max_length and max_length > 0 and len(content) > max_length: 

355 content = content[:max_length] + '\n\n[Content truncated...]' 

356 else: 

357 content = '' 

358 

359 return content 

360 except Exception as e: 

361 log.warning(f'fetch_url error: {e}') 

362 return JSONCodec.dumps({'error': str(e)}) 

363 

364 

365# ============================================================================= 

366# IMAGE GENERATION TOOLS 

367# ============================================================================= 

368 

369 

370async def generate_image( 

371 prompt: str, 

372 __request__: Request = None, 

373 __user__: dict = None, 

374 __event_emitter__: callable = None, 

375 __chat_id__: str = None, 

376 __message_id__: str = None, 

377) -> str: 

378 """ 

379 Generate an image based on a text prompt. 

380 

381 :param prompt: A detailed description of the image to generate 

382 :return: Confirmation that the image was generated, or an error message 

383 """ 

384 if __request__ is None: 

385 return JSONCodec.dumps({'error': 'Request context not available'}) 

386 

387 try: 

388 user = UserModel(**__user__) if __user__ else None 

389 

390 images = await image_generations( 

391 request=__request__, 

392 form_data=CreateImageForm(prompt=prompt), 

393 metadata=( 

394 {'channel_id': __chat_id__.removeprefix('channel:'), 'message_id': __message_id__} 

395 if isinstance(__chat_id__, str) and __chat_id__.startswith('channel:') 

396 else None 

397 ), 

398 user=user, 

399 ) 

400 

401 # Prepare file entries for the images 

402 image_files = [{'type': 'image', **img} for img in images] 

403 

404 # Persist files to DB if chat context is available 

405 if is_saved_chat_id(__chat_id__) and __message_id__ and images: 

406 db_files = await Chats.add_message_files_by_id_and_message_id( 

407 __chat_id__, 

408 __message_id__, 

409 image_files, 

410 ) 

411 if db_files is not None: 

412 image_files = db_files 

413 

414 # Emit the images to the UI if event emitter is available 

415 if __event_emitter__ and image_files: 

416 await __event_emitter__( 

417 { 

418 'type': 'chat:message:files', 

419 'data': { 

420 'files': image_files, 

421 }, 

422 } 

423 ) 

424 # Return a message indicating the image is already displayed 

425 return JSONCodec.dumps( 

426 { 

427 'status': 'success', 

428 'message': 'The image has been successfully generated and is already visible to the user in the chat. You do not need to display or embed the image again - just acknowledge that it has been created.', 

429 'images': images, 

430 }, 

431 ensure_ascii=False, 

432 ) 

433 

434 return JSONCodec.dumps({'status': 'success', 'images': images}, ensure_ascii=False) 

435 except Exception as e: 

436 log.exception(f'generate_image error: {e}') 

437 return JSONCodec.dumps({'error': str(e)}) 

438 

439 

440async def edit_image( 

441 prompt: str, 

442 image_urls: list[str], 

443 __request__: Request = None, 

444 __user__: dict = None, 

445 __event_emitter__: callable = None, 

446 __chat_id__: str = None, 

447 __message_id__: str = None, 

448) -> str: 

449 """ 

450 Transform one or more existing images according to a text prompt. 

451 Supports targeted edits such as adding, removing, replacing, inpainting, extending, or compositing image content. 

452 

453 :param prompt: A description of the transformation to apply to the provided images 

454 :param image_urls: Source image URLs to modify or use as composition inputs 

455 :return: Confirmation that the images were edited, or an error message 

456 """ 

457 if __request__ is None: 

458 return JSONCodec.dumps({'error': 'Request context not available'}) 

459 

460 try: 

461 user = UserModel(**__user__) if __user__ else None 

462 

463 images = await image_edits( 

464 request=__request__, 

465 form_data=EditImageForm(prompt=prompt, image=image_urls), 

466 metadata=( 

467 {'channel_id': __chat_id__.removeprefix('channel:'), 'message_id': __message_id__} 

468 if isinstance(__chat_id__, str) and __chat_id__.startswith('channel:') 

469 else None 

470 ), 

471 user=user, 

472 ) 

473 

474 # Prepare file entries for the images 

475 image_files = [{'type': 'image', **img} for img in images] 

476 

477 # Persist files to DB if chat context is available 

478 if is_saved_chat_id(__chat_id__) and __message_id__ and images: 

479 db_files = await Chats.add_message_files_by_id_and_message_id( 

480 __chat_id__, 

481 __message_id__, 

482 image_files, 

483 ) 

484 if db_files is not None: 

485 image_files = db_files 

486 

487 # Emit the images to the UI if event emitter is available 

488 if __event_emitter__ and image_files: 

489 await __event_emitter__( 

490 { 

491 'type': 'chat:message:files', 

492 'data': { 

493 'files': image_files, 

494 }, 

495 } 

496 ) 

497 # Return a message indicating the image is already displayed 

498 return JSONCodec.dumps( 

499 { 

500 'status': 'success', 

501 'message': 'The edited image has been successfully generated and is already visible to the user in the chat. You do not need to display or embed the image again - just acknowledge that it has been created.', 

502 'images': images, 

503 }, 

504 ensure_ascii=False, 

505 ) 

506 

507 return JSONCodec.dumps({'status': 'success', 'images': images}, ensure_ascii=False) 

508 except Exception as e: 

509 log.exception(f'edit_image error: {e}') 

510 return JSONCodec.dumps({'error': str(e)}) 

511 

512 

513# ============================================================================= 

514# USER INPUT TOOLS 

515# ============================================================================= 

516 

517 

518async def ask_user( 

519 questions: list[dict], 

520 allow_other: bool = True, 

521 timeout_ms: int = 120_000, 

522 __event_call__: callable = None, 

523) -> str: 

524 """ 

525 Ask the user clarifying questions before continuing. 

526 Use this when the next step depends on user intent, preference, or a tradeoff that cannot be inferred safely. 

527 

528 :param questions: 1-3 question objects, each with id, header, question, and 2-3 options. Each option needs label and description. 

529 List the option you recommend first; the UI labels the first option Recommended. 

530 :param allow_other: Whether users may enter a free-form answer instead of choosing one of the options 

531 :param timeout_ms: How long the browser should keep the prompt open before cancelling it 

532 :return: JSON with status and answers keyed by question id 

533 """ 

534 try: 

535 if not isinstance(questions, list) or not 1 <= len(questions) <= 3: 

536 raise ValueError('ask_user requires 1-3 questions.') 

537 

538 normalized_questions = [] 

539 seen_ids = set() 

540 for index, question in enumerate(questions): 

541 if not isinstance(question, dict): 

542 raise ValueError('Each question must be an object.') 

543 

544 question_id = str(question.get('id') or '').strip()[:64] 

545 if not question_id: 

546 raise ValueError('Each question requires a non-empty id.') 

547 if question_id in seen_ids: 

548 raise ValueError(f'Duplicate question id: {question_id}') 

549 seen_ids.add(question_id) 

550 

551 options = question.get('options') 

552 if not isinstance(options, list) or not 2 <= len(options) <= 3: 

553 raise ValueError('Each question requires 2-3 options.') 

554 

555 normalized_options = [] 

556 for option in options: 

557 if not isinstance(option, dict): 

558 raise ValueError('Each option must be an object.') 

559 

560 label = str(option.get('label') or '').strip()[:80] 

561 description = str(option.get('description') or '').strip()[:240] 

562 if not label or not description: 

563 raise ValueError('Each option requires a label and description.') 

564 

565 normalized_options.append( 

566 { 

567 'label': label, 

568 'description': description, 

569 } 

570 ) 

571 

572 question_text = str(question.get('question') or '').strip()[:500] 

573 if not question_text: 

574 raise ValueError('Each question requires question text.') 

575 

576 normalized_questions.append( 

577 { 

578 'id': question_id, 

579 'header': str(question.get('header') or '').strip()[:48] or f'Question {index + 1}', 

580 'question': question_text, 

581 'options': normalized_options, 

582 'allow_other': bool(question.get('allow_other', allow_other)), 

583 } 

584 ) 

585 

586 if isinstance(timeout_ms, bool) or not isinstance(timeout_ms, int) or not 60_000 <= timeout_ms <= 240_000: 

587 timeout_ms = 120_000 

588 

589 if __event_call__ is None: 

590 return JSONCodec.dumps( 

591 { 

592 'status': 'error', 

593 'error': 'User input requires an active browser session with WebSocket connection.', 

594 }, 

595 ensure_ascii=False, 

596 ) 

597 

598 output = await __event_call__( 

599 { 

600 'type': 'request:user_input', 

601 'data': { 

602 'questions': normalized_questions, 

603 'allow_other': allow_other, 

604 'timeout_ms': timeout_ms, 

605 }, 

606 } 

607 ) 

608 

609 if not isinstance(output, dict): 

610 return JSONCodec.dumps({'status': 'error', 'error': 'Invalid user input response.'}, ensure_ascii=False) 

611 if output.get('error'): 

612 return JSONCodec.dumps({'status': 'error', 'error': output.get('error')}, ensure_ascii=False) 

613 if output.get('status') == 'cancelled': 

614 return JSONCodec.dumps({'status': 'cancelled', 'answers': {}}, ensure_ascii=False) 

615 

616 return JSONCodec.dumps( 

617 { 

618 'status': 'answered', 

619 'answers': output.get('answers', {}), 

620 }, 

621 ensure_ascii=False, 

622 ) 

623 except Exception as e: 

624 log.exception(f'ask_user error: {e}') 

625 return JSONCodec.dumps({'status': 'error', 'error': str(e)}, ensure_ascii=False) 

626 

627 

628# ============================================================================= 

629# CODE INTERPRETER TOOLS 

630# ============================================================================= 

631 

632 

633async def execute_code( 

634 code: str, 

635 __request__: Request = None, 

636 __user__: dict = None, 

637 __event_emitter__: callable = None, 

638 __event_call__: callable = None, 

639 __chat_id__: str = None, 

640 __message_id__: str = None, 

641 __metadata__: dict = None, 

642) -> str: 

643 """ 

644 Execute Python code in a sandboxed environment and return the output. 

645 Use this to perform calculations, data analysis, generate visualizations, 

646 or run any Python code that would help answer the user's question. 

647 

648 :param code: The Python code to execute 

649 :return: JSON with stdout, stderr, and result from execution 

650 """ 

651 from uuid import uuid4 

652 

653 if __request__ is None: 

654 return JSONCodec.dumps({'error': 'Request context not available'}) 

655 

656 try: 

657 # Sanitize code (strips ANSI codes and markdown fences) 

658 code = sanitize_code(code) 

659 

660 # Import blocked modules from config (same as middleware) 

661 from open_webui.config import CODE_INTERPRETER_BLOCKED_MODULES 

662 

663 # Add import blocking code if there are blocked modules 

664 if CODE_INTERPRETER_BLOCKED_MODULES: 

665 import textwrap 

666 

667 blocking_code = textwrap.dedent( 

668 f""" 

669 import builtins 

670 

671 BLOCKED_MODULES = {CODE_INTERPRETER_BLOCKED_MODULES} 

672 

673 _real_import = builtins.__import__ 

674 def restricted_import(name, globals=None, locals=None, fromlist=(), level=0): 

675 if name.split('.')[0] in BLOCKED_MODULES: 

676 importer_name = globals.get('__name__') if globals else None 

677 if importer_name == '__main__': 

678 raise ImportError( 

679 f"Direct import of module {{name}} is restricted." 

680 ) 

681 return _real_import(name, globals, locals, fromlist, level) 

682 

683 builtins.__import__ = restricted_import 

684 """ 

685 ) 

686 code = blocking_code + '\n' + code 

687 

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

689 if engine == 'pyodide': 

690 # Execute via frontend pyodide using bidirectional event call 

691 if __event_call__ is None: 

692 return JSONCodec.dumps( 

693 {'error': 'Event call not available. WebSocket connection required for pyodide execution.'} 

694 ) 

695 

696 output = await __event_call__( 

697 { 

698 'type': 'execute:python', 

699 'data': { 

700 'id': str(uuid4()), 

701 'code': code, 

702 'session_id': (__metadata__.get('session_id') if __metadata__ else None), 

703 'files': (__metadata__.get('files', []) if __metadata__ else []), 

704 }, 

705 } 

706 ) 

707 

708 # Parse the output - pyodide returns dict with stdout, stderr, result 

709 if isinstance(output, dict): 

710 # Handle error responses from event_caller (e.g. session disconnected, timeout) 

711 if output.get('error') and not output.get('stdout') and not output.get('result'): 

712 stderr = output['error'] 

713 stdout = '' 

714 result = '' 

715 else: 

716 stdout = output.get('stdout', '') 

717 stderr = output.get('stderr', '') 

718 result = output.get('result', '') 

719 else: 

720 stdout = '' 

721 stderr = '' 

722 result = str(output) if output else '' 

723 

724 elif engine == 'jupyter': 

725 from open_webui.utils.code_interpreter import execute_code_jupyter 

726 

727 jupyter_auth = await Config.get('code_interpreter.jupyter.auth') 

728 

729 output = await execute_code_jupyter( 

730 await Config.get('code_interpreter.jupyter.url'), 

731 code, 

732 (await Config.get('code_interpreter.jupyter.auth_token') if jupyter_auth == 'token' else None), 

733 (await Config.get('code_interpreter.jupyter.auth_password') if jupyter_auth == 'password' else None), 

734 await Config.get('code_interpreter.jupyter.timeout'), 

735 ) 

736 

737 stdout = output.get('stdout', '') 

738 stderr = output.get('stderr', '') 

739 result = output.get('result', '') 

740 

741 else: 

742 return JSONCodec.dumps({'error': f'Unknown code interpreter engine: {engine}'}) 

743 

744 # Handle image outputs (base64 encoded) - replace with uploaded URLs 

745 # Get actual user object for image upload (upload_image requires user.id attribute) 

746 if __user__ and __user__.get('id'): 

747 from open_webui.models.users import Users 

748 from open_webui.utils.files import get_image_url_from_base64 

749 

750 user = await Users.get_user_by_id(__user__['id']) 

751 

752 # Extract and upload images from stdout 

753 if stdout and isinstance(stdout, str): 

754 stdout_lines = stdout.split('\n') 

755 for idx, line in enumerate(stdout_lines): 

756 if 'data:image/png;base64' in line: 

757 image_url = await get_image_url_from_base64( 

758 __request__, 

759 line, 

760 __metadata__ or {}, 

761 user, 

762 ) 

763 if image_url: 

764 stdout_lines[idx] = f'![Output Image]({image_url})' 

765 stdout = '\n'.join(stdout_lines) 

766 

767 # Extract and upload images from result 

768 if result and isinstance(result, str): 

769 result_lines = result.split('\n') 

770 for idx, line in enumerate(result_lines): 

771 if 'data:image/png;base64' in line: 

772 image_url = await get_image_url_from_base64( 

773 __request__, 

774 line, 

775 __metadata__ or {}, 

776 user, 

777 ) 

778 if image_url: 

779 result_lines[idx] = f'![Output Image]({image_url})' 

780 result = '\n'.join(result_lines) 

781 

782 response = { 

783 'status': 'success', 

784 'stdout': stdout, 

785 'stderr': stderr, 

786 'result': result, 

787 } 

788 

789 return JSONCodec.dumps(response, ensure_ascii=False) 

790 except Exception as e: 

791 log.exception(f'execute_code error: {e}') 

792 return JSONCodec.dumps({'error': str(e)}) 

793 

794 

795# ============================================================================= 

796# MEMORY TOOLS 

797# ============================================================================= 

798 

799 

800async def list_memory_paths( 

801 query: str = '', 

802 count: int = 100, 

803 type: str = 'all', 

804 __request__: Request = None, 

805 __user__: dict = None, 

806) -> str: 

807 """ 

808 List saved memory paths to find existing memory groups before writing or moving memories. 

809 

810 :param query: Optional query to filter memory paths or contents 

811 :param count: Maximum number of paths to return 

812 :param type: "user", "context", or "all" 

813 :return: JSON with memory paths, counts, children, and update times 

814 """ 

815 try: 

816 user = UserModel(**__user__) if __user__ else None 

817 result = await _list_memory_paths( 

818 ListMemoryPathsForm( 

819 query=query or None, 

820 type=type if type in {'user', 'context', 'all'} else 'all', 

821 limit=count, 

822 ), 

823 user, 

824 ) 

825 return JSONCodec.dumps(result, ensure_ascii=False) 

826 except Exception as e: 

827 log.exception(f'list_memory_paths error: {e}') 

828 return JSONCodec.dumps({'error': str(e)}) 

829 

830 

831async def read_memory_path( 

832 path: str, 

833 count: int = 50, 

834 type: str = 'all', 

835 include_children: bool = True, 

836 __request__: Request = None, 

837 __user__: dict = None, 

838) -> str: 

839 """ 

840 Read saved memories at a memory path, including nearby parent and child paths. 

841 

842 :param path: Memory path to read 

843 :param count: Maximum number of memories to return 

844 :param type: "user", "context", or "all" 

845 :param include_children: Include memories under child paths 

846 :return: JSON with parent paths, child paths, and memories at the path 

847 """ 

848 try: 

849 user = UserModel(**__user__) if __user__ else None 

850 result = await _read_memory_path( 

851 ReadMemoryPathForm( 

852 path=path, 

853 type=type if type in {'user', 'context', 'all'} else 'all', 

854 include_children=include_children, 

855 limit=count, 

856 ), 

857 user, 

858 ) 

859 return JSONCodec.dumps(result, ensure_ascii=False) 

860 except Exception as e: 

861 log.exception(f'read_memory_path error: {e}') 

862 return JSONCodec.dumps({'error': str(e)}) 

863 

864 

865async def search_memories( 

866 query: str = '', 

867 count: int = 5, 

868 type: str = 'all', 

869 path: Optional[str] = None, 

870 memory_id: Optional[str] = None, 

871 __request__: Request = None, 

872 __user__: dict = None, 

873) -> str: 

874 """ 

875 Search or browse saved memories by content, path, type, or memory ID. 

876 

877 :param query: Optional query to search memory content and path 

878 :param count: Number of memories to return (default 5) 

879 :param type: "user", "context", or "all" 

880 :param path: Optional memory path to search around 

881 :param memory_id: Optional exact memory ID to read 

882 :return: JSON with matching memories and their dates 

883 """ 

884 if __request__ is None: 

885 return JSONCodec.dumps({'error': 'Request context not available'}) 

886 

887 try: 

888 user = UserModel(**__user__) if __user__ else None 

889 

890 memories = await _search_memories( 

891 SearchMemoriesForm( 

892 query=query or None, 

893 type=type if type in {'user', 'context', 'all'} else 'all', 

894 path=path, 

895 memory_id=memory_id, 

896 limit=count, 

897 ), 

898 user, 

899 ) 

900 

901 if not memories: 

902 return JSONCodec.dumps([]) 

903 

904 return JSONCodec.dumps( 

905 [ 

906 { 

907 'id': memory.id, 

908 'type': memory.type, 

909 'path': memory.path, 

910 'content': memory.content, 

911 'created_at': time.strftime('%Y-%m-%d', time.localtime(memory.created_at)), 

912 'updated_at': time.strftime('%Y-%m-%d', time.localtime(memory.updated_at)), 

913 } 

914 for memory in memories 

915 ], 

916 ensure_ascii=False, 

917 ) 

918 except Exception as e: 

919 log.exception(f'search_memories error: {e}') 

920 return JSONCodec.dumps({'error': str(e)}) 

921 

922 

923async def add_memory( 

924 content: str, 

925 type: str = 'user', 

926 path: Optional[str] = None, 

927 __request__: Request = None, 

928 __user__: dict = None, 

929) -> str: 

930 """ 

931 Save enduring information that can improve future chats. 

932 

933 Save stable preferences, goals, projects, relationships, habits, and standing instructions. 

934 Do not save one-off activity, meals, routine daily events, temporary mood, or other short-lived details 

935 unless the user explicitly asks you to remember them. 

936 

937 :param content: The memory content to store 

938 :param type: Use "user" for facts/preferences about the user, or "context" for other durable context 

939 :param path: Optional stable memory address for grouping related memories 

940 :return: Confirmation that the memory was stored 

941 """ 

942 if __request__ is None: 

943 return JSONCodec.dumps({'error': 'Request context not available'}) 

944 

945 try: 

946 user = UserModel(**__user__) if __user__ else None 

947 

948 memory = await _add_memory( 

949 __request__, 

950 AddMemoryForm(content=content, type=Memories.normalize_memory_type(type), path=path), 

951 user, 

952 ) 

953 

954 return JSONCodec.dumps( 

955 {'status': 'success', 'id': memory.id, 'type': memory.type, 'path': memory.path}, 

956 ensure_ascii=False, 

957 ) 

958 except Exception as e: 

959 log.exception(f'add_memory error: {e}') 

960 return JSONCodec.dumps({'error': str(e)}) 

961 

962 

963async def update_memory( 

964 operations: list[dict], 

965 __request__: Request = None, 

966 __user__: dict = None, 

967) -> str: 

968 """ 

969 Apply a batch of memory changes after learning enduring information. 

970 

971 Use type "user" for facts, preferences, or instructions about the user. 

972 Use type "context" for other durable context that may help future chats. 

973 Do not save one-off activity, meals, routine daily events, temporary mood, or other short-lived details 

974 unless the user explicitly asks you to remember them. 

975 Path is optional. Use it as a stable memory address to group related memories. 

976 Prefer an existing path from list_memory_paths when one fits. 

977 Leave path empty when no useful grouping is clear. 

978 

979 Operation shapes: 

980 - {"action": "add", "content": "...", "type": "user"|"context", "path": "..."} 

981 - {"action": "replace", "id": "...", "content": "...", "type": "user"|"context", "path": "..."} 

982 - {"action": "move", "id": "...", "path": "..."} 

983 - {"action": "remove", "id": "..."} 

984 

985 :param operations: Memory operations to apply in one request 

986 :return: JSON with operation results 

987 """ 

988 if __request__ is None: 

989 return JSONCodec.dumps({'error': 'Request context not available'}) 

990 

991 try: 

992 user = UserModel(**__user__) if __user__ else None 

993 operation_results = await _update_memories( 

994 __request__, 

995 UpdateMemoriesForm(operations=operations), 

996 user, 

997 ) 

998 return JSONCodec.dumps(operation_results, ensure_ascii=False) 

999 except Exception as e: 

1000 log.exception(f'update_memory error: {e}') 

1001 return JSONCodec.dumps({'error': str(e)}) 

1002 

1003 

1004async def replace_memory_content( 

1005 memory_id: str, 

1006 content: str, 

1007 type: Optional[str] = None, 

1008 path: Optional[str] = None, 

1009 __request__: Request = None, 

1010 __user__: dict = None, 

1011) -> str: 

1012 """ 

1013 Update an existing saved memory by its ID when its content needs correction. 

1014 

1015 :param memory_id: The ID of the memory to update 

1016 :param content: The new content for the memory 

1017 :param type: Optional "user" or "context" type for the updated memory 

1018 :param path: Optional stable memory address for grouping related memories 

1019 :return: Confirmation that the memory was updated 

1020 """ 

1021 if __request__ is None: 

1022 return JSONCodec.dumps({'error': 'Request context not available'}) 

1023 

1024 try: 

1025 user = UserModel(**__user__) if __user__ else None 

1026 

1027 memory = await update_memory_by_id( 

1028 memory_id=memory_id, 

1029 request=__request__, 

1030 form_data=MemoryUpdateModel( 

1031 content=content, 

1032 type=Memories.normalize_memory_type(type) if type else None, 

1033 path=path, 

1034 ), 

1035 user=user, 

1036 ) 

1037 

1038 return JSONCodec.dumps( 

1039 { 

1040 'status': 'success', 

1041 'id': memory.id, 

1042 'type': memory.type, 

1043 'path': memory.path, 

1044 'content': memory.content, 

1045 }, 

1046 ensure_ascii=False, 

1047 ) 

1048 except Exception as e: 

1049 log.exception(f'replace_memory_content error: {e}') 

1050 return JSONCodec.dumps({'error': str(e)}) 

1051 

1052 

1053async def delete_memory( 

1054 memory_id: str, 

1055 __request__: Request = None, 

1056 __user__: dict = None, 

1057) -> str: 

1058 """ 

1059 Delete a saved memory by its ID. 

1060 

1061 :param memory_id: The ID of the memory to delete 

1062 :return: Confirmation that the memory was deleted 

1063 """ 

1064 if __request__ is None: 

1065 return JSONCodec.dumps({'error': 'Request context not available'}) 

1066 

1067 try: 

1068 user = UserModel(**__user__) if __user__ else None 

1069 

1070 result = await Memories.delete_memory_by_id_and_user_id(memory_id, user.id) 

1071 

1072 if result: 

1073 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=f'user-memory-{user.id}', ids=[memory_id]) 

1074 return JSONCodec.dumps( 

1075 {'status': 'success', 'message': f'Memory {memory_id} deleted'}, 

1076 ensure_ascii=False, 

1077 ) 

1078 else: 

1079 return JSONCodec.dumps({'error': 'Memory not found or access denied'}) 

1080 except Exception as e: 

1081 log.exception(f'delete_memory error: {e}') 

1082 return JSONCodec.dumps({'error': str(e)}) 

1083 

1084 

1085async def list_memories( 

1086 __request__: Request = None, 

1087 __user__: dict = None, 

1088) -> str: 

1089 """ 

1090 List all stored memories for the user, including IDs and timestamps. 

1091 

1092 :return: JSON list of all memories with id, content, and dates 

1093 """ 

1094 if __request__ is None: 

1095 return JSONCodec.dumps({'error': 'Request context not available'}) 

1096 

1097 try: 

1098 user = UserModel(**__user__) if __user__ else None 

1099 

1100 memories = await Memories.get_memories_by_user_id(user.id) 

1101 

1102 if memories: 

1103 memory_rows = [ 

1104 { 

1105 'id': m.id, 

1106 'type': m.type, 

1107 'path': m.path, 

1108 'content': m.content, 

1109 'created_at': time.strftime('%Y-%m-%d %H:%M', time.localtime(m.created_at)), 

1110 'updated_at': time.strftime('%Y-%m-%d %H:%M', time.localtime(m.updated_at)), 

1111 } 

1112 for m in memories 

1113 ] 

1114 return JSONCodec.dumps(memory_rows, ensure_ascii=False) 

1115 else: 

1116 return JSONCodec.dumps([]) 

1117 except Exception as e: 

1118 log.exception(f'list_memories error: {e}') 

1119 return JSONCodec.dumps({'error': str(e)}) 

1120 

1121 

1122# ============================================================================= 

1123# NOTES TOOLS 

1124# ============================================================================= 

1125 

1126 

1127async def search_notes( 

1128 query: str, 

1129 count: int = 5, 

1130 start_timestamp: Optional[int] = None, 

1131 end_timestamp: Optional[int] = None, 

1132 __request__: Request = None, 

1133 __user__: dict = None, 

1134) -> str: 

1135 """ 

1136 Search the user's saved notes by title and content. 

1137 

1138 :param query: The search query to find matching notes 

1139 :param count: Maximum number of results to return (default: 5) 

1140 :param start_timestamp: Only include notes updated after this Unix timestamp (seconds) 

1141 :param end_timestamp: Only include notes updated before this Unix timestamp (seconds) 

1142 :return: JSON with matching notes containing id, title, and content snippet 

1143 """ 

1144 if __request__ is None: 

1145 return JSONCodec.dumps({'error': 'Request context not available'}) 

1146 

1147 if not __user__: 

1148 return JSONCodec.dumps({'error': 'User context not available'}) 

1149 

1150 try: 

1151 user_id = __user__.get('id') 

1152 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

1153 

1154 result = await Notes.search_notes( 

1155 user_id=user_id, 

1156 filter={ 

1157 'query': query, 

1158 'user_id': user_id, 

1159 'group_ids': user_group_ids, 

1160 'permission': 'read', 

1161 }, 

1162 skip=0, 

1163 limit=count * 3, # Fetch more for filtering 

1164 ) 

1165 

1166 # Convert timestamps to nanoseconds for comparison 

1167 start_ts = start_timestamp * 1_000_000_000 if start_timestamp else None 

1168 end_ts = end_timestamp * 1_000_000_000 if end_timestamp else None 

1169 

1170 notes = [] 

1171 for note in result.items: 

1172 # Apply date filters (updated_at is in nanoseconds) 

1173 if start_ts and note.updated_at < start_ts: 

1174 continue 

1175 if end_ts and note.updated_at > end_ts: 

1176 continue 

1177 

1178 # Extract a snippet from the markdown content 

1179 content_snippet = '' 

1180 if note.data and note.data.get('content', {}).get('md'): 

1181 md_content = note.data['content']['md'] 

1182 content_lower = md_content.lower() 

1183 

1184 # Find the first matching word to center the snippet around. 

1185 search_words = query.lower().split() 

1186 match_pos = -1 

1187 match_len = len(query) 

1188 for word in search_words: 

1189 found_pos = content_lower.find(word) 

1190 if found_pos != -1: 

1191 match_pos = found_pos 

1192 match_len = len(word) 

1193 break 

1194 

1195 if match_pos != -1: 

1196 snippet_start = max(0, match_pos - 50) 

1197 snippet_end = min(len(md_content), match_pos + match_len + 100) 

1198 content_snippet = ( 

1199 ('...' if snippet_start > 0 else '') 

1200 + md_content[snippet_start:snippet_end] 

1201 + ('...' if snippet_end < len(md_content) else '') 

1202 ) 

1203 else: 

1204 content_snippet = md_content[:150] + ('...' if len(md_content) > 150 else '') 

1205 

1206 notes.append( 

1207 { 

1208 'id': note.id, 

1209 'title': note.title, 

1210 'snippet': content_snippet, 

1211 'updated_at': note.updated_at, 

1212 } 

1213 ) 

1214 

1215 if len(notes) >= count: 

1216 break 

1217 

1218 return JSONCodec.dumps(notes, ensure_ascii=False) 

1219 except Exception as e: 

1220 log.exception(f'search_notes error: {e}') 

1221 return JSONCodec.dumps({'error': str(e)}) 

1222 

1223 

1224async def view_note( 

1225 note_id: str, 

1226 __request__: Request = None, 

1227 __user__: dict = None, 

1228) -> str: 

1229 """ 

1230 Get the full content of a note by its ID. 

1231 

1232 :param note_id: The ID of the note to retrieve 

1233 :return: JSON with the note's id, title, and full markdown content 

1234 """ 

1235 if __request__ is None: 

1236 return JSONCodec.dumps({'error': 'Request context not available'}) 

1237 

1238 if not __user__: 

1239 return JSONCodec.dumps({'error': 'User context not available'}) 

1240 

1241 try: 

1242 note = await Notes.get_note_by_id(note_id) 

1243 

1244 if not note: 

1245 return JSONCodec.dumps({'error': 'Note not found'}) 

1246 

1247 # Check access permission 

1248 user_id = __user__.get('id') 

1249 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

1250 

1251 from open_webui.models.access_grants import AccessGrants 

1252 

1253 if ( 

1254 __user__.get('role') != 'admin' 

1255 and note.user_id != user_id 

1256 and not await AccessGrants.has_access( 

1257 user_id=user_id, 

1258 resource_type='note', 

1259 resource_id=note.id, 

1260 permission='read', 

1261 user_group_ids=set(user_group_ids), 

1262 ) 

1263 ): 

1264 return JSONCodec.dumps({'error': 'Access denied'}) 

1265 

1266 # Extract markdown content 

1267 content = '' 

1268 if note.data and note.data.get('content', {}).get('md'): 

1269 content = note.data['content']['md'] 

1270 

1271 return JSONCodec.dumps( 

1272 { 

1273 'id': note.id, 

1274 'title': note.title, 

1275 'content': content, 

1276 'updated_at': note.updated_at, 

1277 'created_at': note.created_at, 

1278 }, 

1279 ensure_ascii=False, 

1280 ) 

1281 except Exception as e: 

1282 log.exception(f'view_note error: {e}') 

1283 return JSONCodec.dumps({'error': str(e)}) 

1284 

1285 

1286async def write_note( 

1287 title: str, 

1288 content: str, 

1289 __request__: Request = None, 

1290 __user__: dict = None, 

1291) -> str: 

1292 """ 

1293 Create a new note with the given title and content. 

1294 

1295 :param title: The title of the new note 

1296 :param content: The markdown content for the note 

1297 :return: JSON with success status and new note id 

1298 """ 

1299 if __request__ is None: 

1300 return JSONCodec.dumps({'error': 'Request context not available'}) 

1301 

1302 if not __user__: 

1303 return JSONCodec.dumps({'error': 'User context not available'}) 

1304 

1305 try: 

1306 from open_webui.models.notes import NoteForm 

1307 

1308 user_id = __user__.get('id') 

1309 

1310 form = NoteForm( 

1311 title=title, 

1312 data={'content': {'md': content}}, 

1313 access_grants=[], # Private by default - only owner can access 

1314 ) 

1315 

1316 new_note = await Notes.insert_new_note(user_id, form) 

1317 

1318 if not new_note: 

1319 return JSONCodec.dumps({'error': 'Failed to create note'}) 

1320 

1321 return JSONCodec.dumps( 

1322 { 

1323 'status': 'success', 

1324 'id': new_note.id, 

1325 'title': new_note.title, 

1326 'created_at': new_note.created_at, 

1327 }, 

1328 ensure_ascii=False, 

1329 ) 

1330 except Exception as e: 

1331 log.exception(f'write_note error: {e}') 

1332 return JSONCodec.dumps({'error': str(e)}) 

1333 

1334 

1335async def replace_note_content( 

1336 note_id: str, 

1337 content: Optional[str] = None, 

1338 operations: Optional[list[dict]] = None, 

1339 title: Optional[str] = None, 

1340 __request__: Request = None, 

1341 __user__: dict = None, 

1342) -> str: 

1343 """ 

1344 Update an existing note by replacing the whole markdown content or applying range operations. 

1345 

1346 Prefer "replace_range" when only part of the note changes. 

1347 A "replace" operation must be the only operation in the request. 

1348 start and end are 0-indexed character offsets into the markdown content from view_note. 

1349 end is exclusive. 

1350 Offsets never shift as operations are applied. 

1351 Ranges must not overlap. 

1352 expected is optional. When set, the request is rejected if the range's current text does not match it. 

1353 

1354 :param note_id: The ID of the note to update 

1355 :param content: The new markdown content for a whole-note update 

1356 :param operations: Optional note operations: 

1357 - {"action": "replace", "content": "..."} 

1358 - {"action": "replace_range", "start": 0, "end": 10, "content": "...", "expected": "..."} 

1359 :param title: Optional new title for the note 

1360 :return: JSON with success status and updated note info 

1361 """ 

1362 if __request__ is None: 

1363 return JSONCodec.dumps({'error': 'Request context not available'}) 

1364 

1365 if not __user__: 

1366 return JSONCodec.dumps({'error': 'User context not available'}) 

1367 

1368 try: 

1369 from open_webui.models.notes import NoteUpdateForm 

1370 

1371 note = await Notes.get_note_by_id(note_id) 

1372 

1373 if not note: 

1374 return JSONCodec.dumps({'error': 'Note not found', 'code': 'not_found'}) 

1375 

1376 user_id = __user__.get('id') 

1377 if __user__.get('role') != 'admin' and not await _has_write_access_to_note(note, user_id): 

1378 return JSONCodec.dumps({'error': 'Write access denied', 'code': 'write_access_denied'}) 

1379 

1380 current_content = ((note.data or {}).get('content') or {}).get('md') or '' 

1381 applied_operation_count = 0 

1382 if operations is not None: 

1383 if not isinstance(operations, list) or len(operations) == 0: 

1384 return JSONCodec.dumps({'error': 'operations must be a non-empty list', 'code': 'invalid_operations'}) 

1385 

1386 range_operations = [] 

1387 for idx, operation in enumerate(operations): 

1388 if not isinstance(operation, dict): 

1389 return JSONCodec.dumps( 

1390 {'error': 'each operation must be an object', 'code': 'invalid_operation', 'index': idx} 

1391 ) 

1392 

1393 action = operation.get('action') 

1394 replacement = operation.get('content') 

1395 

1396 if action == 'replace': 

1397 if len(operations) != 1: 

1398 return JSONCodec.dumps( 

1399 { 

1400 'error': 'replace operation must be the only operation', 

1401 'code': 'invalid_operations', 

1402 'index': idx, 

1403 } 

1404 ) 

1405 if not isinstance(replacement, str): 

1406 return JSONCodec.dumps( 

1407 { 

1408 'error': 'replace operation content must be a string', 

1409 'code': 'invalid_content', 

1410 'index': idx, 

1411 } 

1412 ) 

1413 content = replacement 

1414 applied_operation_count = 1 

1415 break 

1416 

1417 if action != 'replace_range': 

1418 return JSONCodec.dumps( 

1419 {'error': 'unknown operation action', 'code': 'invalid_action', 'index': idx, 'action': action} 

1420 ) 

1421 

1422 start = operation.get('start') 

1423 end = operation.get('end') 

1424 expected = operation.get('expected') 

1425 if not isinstance(start, int) or not isinstance(end, int): 

1426 return JSONCodec.dumps( 

1427 {'error': 'operation start and end must be integers', 'code': 'invalid_range', 'index': idx} 

1428 ) 

1429 if not isinstance(replacement, str): 

1430 return JSONCodec.dumps( 

1431 {'error': 'operation content must be a string', 'code': 'invalid_content', 'index': idx} 

1432 ) 

1433 if start < 0 or end < start or end > len(current_content): 

1434 return JSONCodec.dumps( 

1435 {'error': 'operation range is out of bounds', 'code': 'range_out_of_bounds', 'index': idx} 

1436 ) 

1437 if expected is not None and current_content[start:end] != expected: 

1438 return JSONCodec.dumps( 

1439 { 

1440 'error': 'operation expected text does not match current content', 

1441 'code': 'expected_mismatch', 

1442 'index': idx, 

1443 } 

1444 ) 

1445 

1446 range_operations.append({'start': start, 'end': end, 'content': replacement}) 

1447 

1448 range_operations.sort(key=lambda operation: operation['start']) 

1449 previous_end = 0 

1450 for idx, operation in enumerate(range_operations): 

1451 if operation['start'] < previous_end: 

1452 return JSONCodec.dumps( 

1453 {'error': 'operation ranges must not overlap', 'code': 'overlapping_operations', 'index': idx} 

1454 ) 

1455 previous_end = operation['end'] 

1456 

1457 if range_operations: 

1458 content = current_content 

1459 for operation in reversed(range_operations): 

1460 content = content[: operation['start']] + operation['content'] + content[operation['end'] :] 

1461 applied_operation_count = len(range_operations) 

1462 elif content is None: 

1463 return JSONCodec.dumps({'error': 'content or operations is required', 'code': 'content_required'}) 

1464 

1465 try: 

1466 await stop_item_tasks(__request__.app.state.redis, f'note:{note_id}') 

1467 except Exception: 

1468 pass 

1469 

1470 update_data = { 

1471 'data': { 

1472 **(note.data or {}), 

1473 'content': { 

1474 **((note.data or {}).get('content') or {}), 

1475 'json': None, 

1476 'html': '', 

1477 'md': content, 

1478 }, 

1479 } 

1480 } 

1481 if title: 

1482 update_data['title'] = title 

1483 

1484 form = NoteUpdateForm(**update_data) 

1485 updated_note = await Notes.update_note_by_id(note_id, form) 

1486 

1487 if not updated_note: 

1488 return JSONCodec.dumps({'error': 'Failed to update note', 'code': 'update_failed'}) 

1489 

1490 await _emit_note_updated(__request__, __user__, updated_note) 

1491 

1492 return JSONCodec.dumps( 

1493 { 

1494 'status': 'success', 

1495 'id': updated_note.id, 

1496 'title': updated_note.title, 

1497 'updated_at': updated_note.updated_at, 

1498 'applied_operation_count': applied_operation_count, 

1499 }, 

1500 ensure_ascii=False, 

1501 ) 

1502 except Exception as e: 

1503 log.exception(f'replace_note_content error: {e}') 

1504 return JSONCodec.dumps({'error': str(e), 'code': 'unexpected_error'}) 

1505 

1506 

1507# ============================================================================= 

1508# CHATS TOOLS 

1509# ============================================================================= 

1510 

1511 

1512async def search_chats( 

1513 query: str, 

1514 count: int = 5, 

1515 start_timestamp: Optional[int] = None, 

1516 end_timestamp: Optional[int] = None, 

1517 __request__: Request = None, 

1518 __user__: dict = None, 

1519 __chat_id__: str = None, 

1520) -> str: 

1521 """ 

1522 Search the user's previous chat conversations by title and message content, 

1523 excluding the current chat. Helpful for finding details from earlier 

1524 conversations when they are not already visible in the current context. 

1525 Exact phrase matches are preferred, and descriptive keyword queries are 

1526 supported. 

1527 

1528 :param query: Exact phrase or descriptive keyword query to find matching previous chats 

1529 :param count: Maximum number of results to return (default: 5) 

1530 :param start_timestamp: Only include chats updated after this Unix timestamp (seconds) 

1531 :param end_timestamp: Only include chats updated before this Unix timestamp (seconds) 

1532 :return: JSON with matching chats containing id, title, updated_at, and content snippet 

1533 """ 

1534 if __request__ is None: 

1535 return JSONCodec.dumps({'error': 'Request context not available'}) 

1536 

1537 if not __user__: 

1538 return JSONCodec.dumps({'error': 'User context not available'}) 

1539 

1540 try: 

1541 user_id = __user__.get('id') 

1542 

1543 chats = await Chats.get_chats_by_user_id_and_search_text( 

1544 user_id=user_id, 

1545 search_text=query, 

1546 include_archived=False, 

1547 skip=0, 

1548 limit=count * 3, # Fetch more for filtering 

1549 ) 

1550 

1551 results = [] 

1552 for chat in chats: 

1553 # Skip the current chat to avoid showing it in search results 

1554 if __chat_id__ and chat.id == __chat_id__: 

1555 continue 

1556 

1557 # Apply date filters (updated_at is in seconds) 

1558 if start_timestamp and chat.updated_at < start_timestamp: 

1559 continue 

1560 if end_timestamp and chat.updated_at > end_timestamp: 

1561 continue 

1562 

1563 # Find a matching message snippet 

1564 snippet = '' 

1565 messages = (getattr(chat, 'chat', None) or {}).get('history', {}).get('messages', {}) 

1566 if not messages: 

1567 messages = (getattr(chat, 'chat', None) or {}).get('messages', {}) or {} 

1568 if isinstance(messages, list): 

1569 messages = {str(idx): message for idx, message in enumerate(messages)} 

1570 

1571 lower_query = chat_search_content_query(query) 

1572 needles = list(dict.fromkeys([lower_query, *chat_search_terms(lower_query)])) if lower_query else [] 

1573 

1574 for needle in needles: 

1575 for msg_id, msg in messages.items(): 

1576 content = msg.get('content', '') if isinstance(msg, dict) else '' 

1577 if isinstance(content, str) and needle in content.lower(): 

1578 idx = content.lower().find(needle) 

1579 start = max(0, idx - 50) 

1580 end = min(len(content), idx + len(needle) + 100) 

1581 snippet = ( 

1582 ('...' if start > 0 else '') + content[start:end] + ('...' if end < len(content) else '') 

1583 ) 

1584 break 

1585 if snippet: 

1586 break 

1587 

1588 title = chat.title or '' 

1589 if not snippet and any(needle in title.lower() for needle in needles): 

1590 snippet = f'Title match: {title}' 

1591 

1592 results.append( 

1593 { 

1594 'id': chat.id, 

1595 'title': chat.title, 

1596 'snippet': snippet, 

1597 'updated_at': chat.updated_at, 

1598 } 

1599 ) 

1600 

1601 if len(results) >= count: 

1602 break 

1603 

1604 return JSONCodec.dumps(results, ensure_ascii=False) 

1605 except Exception as e: 

1606 log.exception(f'search_chats error: {e}') 

1607 return JSONCodec.dumps({'error': str(e)}) 

1608 

1609 

1610async def view_chat( 

1611 chat_id: str, 

1612 __request__: Request = None, 

1613 __user__: dict = None, 

1614) -> str: 

1615 """ 

1616 Get the full conversation history of a chat by its ID after a relevant 

1617 previous chat has been identified. 

1618 

1619 :param chat_id: The ID of the chat to retrieve 

1620 :return: JSON with the chat's id, title, and messages 

1621 """ 

1622 if __request__ is None: 

1623 return JSONCodec.dumps({'error': 'Request context not available'}) 

1624 

1625 if not __user__: 

1626 return JSONCodec.dumps({'error': 'User context not available'}) 

1627 

1628 try: 

1629 user_id = __user__.get('id') 

1630 

1631 chat = await Chats.get_chat_by_id_and_user_id(chat_id, user_id) 

1632 

1633 if not chat: 

1634 return JSONCodec.dumps({'error': 'Chat not found or access denied'}) 

1635 

1636 # Extract messages from history 

1637 messages = [] 

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

1639 msg_dict = history.get('messages', {}) 

1640 

1641 # Build message chain from currentId 

1642 current_id = history.get('currentId') 

1643 visited = set() 

1644 

1645 while current_id and current_id not in visited: 

1646 visited.add(current_id) 

1647 msg = msg_dict.get(current_id) 

1648 if msg: 

1649 messages.append( 

1650 { 

1651 'role': msg.get('role', ''), 

1652 'content': msg.get('content', ''), 

1653 } 

1654 ) 

1655 current_id = msg.get('parentId') if msg else None 

1656 

1657 # Reverse to get chronological order 

1658 messages.reverse() 

1659 

1660 return JSONCodec.dumps( 

1661 { 

1662 'id': chat.id, 

1663 'title': chat.title, 

1664 'messages': messages, 

1665 'updated_at': chat.updated_at, 

1666 'created_at': chat.created_at, 

1667 }, 

1668 ensure_ascii=False, 

1669 ) 

1670 except Exception as e: 

1671 log.exception(f'view_chat error: {e}') 

1672 return JSONCodec.dumps({'error': str(e)}) 

1673 

1674 

1675# ============================================================================= 

1676# SUB-AGENT TOOL 

1677# ============================================================================= 

1678 

1679 

1680async def delegate_task( 

1681 task: str, 

1682 context: str = '', 

1683 file_ids: list[str] | None = None, 

1684 background: bool = False, 

1685 __request__: Request = None, 

1686 __user__: dict = None, 

1687 __metadata__: dict = None, 

1688 __chat_id__: str = None, 

1689 __message_id__: str = None, 

1690) -> str: 

1691 """ 

1692 Delegate focused work to a parallel sub-agent using the current model and tools. 

1693 

1694 :param task: The specific task for the sub-agent to complete 

1695 :param context: Relevant context, decisions, or file paths for the task 

1696 :param file_ids: Attached file IDs the sub-agent needs. Use this for images or files; 

1697 do not put file IDs only in context. 

1698 :param background: Return immediately and continue this chat when the sub-agent finishes 

1699 :return: Foreground result text, or a JSON dispatch handle for background work 

1700 """ 

1701 if __request__ is None: 

1702 return 'Error: request context not available.' 

1703 if getattr(__request__.state, 'internal', False) is True: 

1704 return 'Error: sub-agents cannot delegate recursively.' 

1705 

1706 from open_webui.utils.subagents import delegate 

1707 

1708 return await delegate( 

1709 task, 

1710 context, 

1711 background, 

1712 file_ids=file_ids, 

1713 request=__request__, 

1714 user_data=__user__ or {}, 

1715 metadata=__metadata__ or {}, 

1716 parent_chat_id=__chat_id__ or '', 

1717 parent_message_id=__message_id__, 

1718 ) 

1719 

1720 

1721async def timer( 

1722 prompt: str, 

1723 at: str, 

1724 cancel_on: list[Literal['chat.read', 'chat.user_message']] | None = None, 

1725 __request__: Request = None, 

1726 __user__: dict = None, 

1727 __metadata__: dict = None, 

1728 __chat_id__: str = None, 

1729 __message_id__: str = None, 

1730) -> str: 

1731 """ 

1732 Set a one-shot timer for this chat. 

1733 

1734 :param prompt: The prompt to send back into this chat when the timer fires 

1735 :param at: Relative time like 10s, 5m, 1h, 2d, or a timezone-aware RFC 3339 timestamp 

1736 :param cancel_on: Optional events that cancel the timer before it fires 

1737 :return: JSON status with the scheduled time, or an error string 

1738 """ 

1739 if __request__ is None: 

1740 return 'Error: request context not available.' 

1741 if getattr(__request__.state, 'internal', False) is True: 

1742 return 'Error: timers cannot be set from internal chats.' 

1743 

1744 from open_webui.utils.timers import create_timer 

1745 

1746 return await create_timer( 

1747 prompt=prompt, 

1748 at=at, 

1749 cancel_on=cancel_on, 

1750 request=__request__, 

1751 user_data=__user__ or {}, 

1752 metadata=__metadata__ or {}, 

1753 parent_chat_id=__chat_id__ or '', 

1754 parent_message_id=__message_id__, 

1755 ) 

1756 

1757 

1758# ============================================================================= 

1759# CHANNELS TOOLS 

1760# ============================================================================= 

1761 

1762 

1763async def search_channels( 

1764 query: str, 

1765 count: int = 5, 

1766 __request__: Request = None, 

1767 __user__: dict = None, 

1768) -> str: 

1769 """ 

1770 Search channels by name and description to find accessible team spaces. 

1771 

1772 :param query: The search query to find matching channels 

1773 :param count: Maximum number of results to return (default: 5) 

1774 :return: JSON with matching channels containing id, name, description, and type 

1775 """ 

1776 if __request__ is None: 

1777 return JSONCodec.dumps({'error': 'Request context not available'}) 

1778 

1779 if not __user__: 

1780 return JSONCodec.dumps({'error': 'User context not available'}) 

1781 

1782 try: 

1783 user_id = __user__.get('id') 

1784 

1785 # Get all channels the user has access to 

1786 all_channels = await Channels.get_channels_by_user_id(user_id) 

1787 

1788 # Filter by query 

1789 lower_query = query.lower() 

1790 matching_channels = [] 

1791 

1792 for channel in all_channels: 

1793 name_match = lower_query in channel.name.lower() if channel.name else False 

1794 desc_match = lower_query in (channel.description or '').lower() 

1795 

1796 if name_match or desc_match: 

1797 matching_channels.append( 

1798 { 

1799 'id': channel.id, 

1800 'name': channel.name, 

1801 'description': channel.description or '', 

1802 'type': channel.type or 'public', 

1803 } 

1804 ) 

1805 

1806 if len(matching_channels) >= count: 

1807 break 

1808 

1809 return JSONCodec.dumps(matching_channels, ensure_ascii=False) 

1810 except Exception as e: 

1811 log.exception(f'search_channels error: {e}') 

1812 return JSONCodec.dumps({'error': str(e)}) 

1813 

1814 

1815async def search_channel_messages( 

1816 query: str, 

1817 count: int = 10, 

1818 start_timestamp: Optional[int] = None, 

1819 end_timestamp: Optional[int] = None, 

1820 __request__: Request = None, 

1821 __user__: dict = None, 

1822) -> str: 

1823 """ 

1824 Search messages in channels the user is a member of, including thread replies. 

1825 Helpful for finding prior team/channel discussion. 

1826 

1827 :param query: The search query to find matching messages 

1828 :param count: Maximum number of results to return (default: 10) 

1829 :param start_timestamp: Only include messages created after this Unix timestamp (seconds) 

1830 :param end_timestamp: Only include messages created before this Unix timestamp (seconds) 

1831 :return: JSON with matching messages containing channel info, message content, and thread context 

1832 """ 

1833 if __request__ is None: 

1834 return JSONCodec.dumps({'error': 'Request context not available'}) 

1835 

1836 if not __user__: 

1837 return JSONCodec.dumps({'error': 'User context not available'}) 

1838 

1839 try: 

1840 user_id = __user__.get('id') 

1841 

1842 # Get all channels the user has access to 

1843 user_channels = await Channels.get_channels_by_user_id(user_id) 

1844 channel_ids = [c.id for c in user_channels] 

1845 channel_map = {c.id: c for c in user_channels} 

1846 

1847 if not channel_ids: 

1848 return JSONCodec.dumps([]) 

1849 

1850 # Convert timestamps to nanoseconds (Message.created_at is in nanoseconds) 

1851 start_ts = start_timestamp * 1_000_000_000 if start_timestamp else None 

1852 end_ts = end_timestamp * 1_000_000_000 if end_timestamp else None 

1853 

1854 # Search messages using the model method 

1855 matching_messages = await Messages.search_messages_by_channel_ids( 

1856 channel_ids=channel_ids, 

1857 query=query, 

1858 start_timestamp=start_ts, 

1859 end_timestamp=end_ts, 

1860 limit=count, 

1861 ) 

1862 

1863 results = [] 

1864 for msg in matching_messages: 

1865 channel = channel_map.get(msg.channel_id) 

1866 

1867 # Extract snippet around the match 

1868 content = msg.content or '' 

1869 lower_query = query.lower() 

1870 idx = content.lower().find(lower_query) 

1871 if idx != -1: 

1872 start = max(0, idx - 50) 

1873 end = min(len(content), idx + len(query) + 100) 

1874 snippet = ('...' if start > 0 else '') + content[start:end] + ('...' if end < len(content) else '') 

1875 else: 

1876 snippet = content[:150] + ('...' if len(content) > 150 else '') 

1877 

1878 results.append( 

1879 { 

1880 'channel_id': msg.channel_id, 

1881 'channel_name': channel.name if channel else 'Unknown', 

1882 'message_id': msg.id, 

1883 'content_snippet': snippet, 

1884 'is_thread_reply': msg.parent_id is not None, 

1885 'parent_id': msg.parent_id, 

1886 'created_at': msg.created_at, 

1887 } 

1888 ) 

1889 

1890 return JSONCodec.dumps(results, ensure_ascii=False) 

1891 except Exception as e: 

1892 log.exception(f'search_channel_messages error: {e}') 

1893 return JSONCodec.dumps({'error': str(e)}) 

1894 

1895 

1896async def view_channel_message( 

1897 message_id: str, 

1898 __request__: Request = None, 

1899 __user__: dict = None, 

1900) -> str: 

1901 """ 

1902 Get the full content of a channel message by its ID, including thread replies. 

1903 

1904 :param message_id: The ID of the message to retrieve 

1905 :return: JSON with the message content, channel info, and thread replies if any 

1906 """ 

1907 if __request__ is None: 

1908 return JSONCodec.dumps({'error': 'Request context not available'}) 

1909 

1910 if not __user__: 

1911 return JSONCodec.dumps({'error': 'User context not available'}) 

1912 

1913 try: 

1914 user_id = __user__.get('id') 

1915 

1916 message = await Messages.get_message_by_id(message_id) 

1917 

1918 if not message: 

1919 return JSONCodec.dumps({'error': 'Message not found'}) 

1920 

1921 # Verify user has access to the channel 

1922 channel = await Channels.get_channel_by_id(message.channel_id) 

1923 if not channel: 

1924 return JSONCodec.dumps({'error': 'Channel not found'}) 

1925 

1926 # Check if user has access to the channel 

1927 user_channels = await Channels.get_channels_by_user_id(user_id) 

1928 channel_ids = [c.id for c in user_channels] 

1929 

1930 if message.channel_id not in channel_ids: 

1931 return JSONCodec.dumps({'error': 'Access denied'}) 

1932 

1933 # Build response with thread information 

1934 result = { 

1935 'id': message.id, 

1936 'channel_id': message.channel_id, 

1937 'channel_name': channel.name, 

1938 'content': message.content, 

1939 'user_id': message.user_id, 

1940 'is_thread_reply': message.parent_id is not None, 

1941 'parent_id': message.parent_id, 

1942 'reply_count': message.reply_count, 

1943 'created_at': message.created_at, 

1944 'updated_at': message.updated_at, 

1945 } 

1946 

1947 # Include user info if available 

1948 if message.user: 

1949 result['user_name'] = message.user.name 

1950 

1951 return JSONCodec.dumps(result, ensure_ascii=False) 

1952 except Exception as e: 

1953 log.exception(f'view_channel_message error: {e}') 

1954 return JSONCodec.dumps({'error': str(e)}) 

1955 

1956 

1957async def view_channel_thread( 

1958 parent_message_id: str, 

1959 __request__: Request = None, 

1960 __user__: dict = None, 

1961) -> str: 

1962 """ 

1963 Get all messages in a channel thread, including the parent message and all replies. 

1964 

1965 :param parent_message_id: The ID of the parent message that started the thread 

1966 :return: JSON with the parent message and all thread replies in chronological order 

1967 """ 

1968 if __request__ is None: 

1969 return JSONCodec.dumps({'error': 'Request context not available'}) 

1970 

1971 if not __user__: 

1972 return JSONCodec.dumps({'error': 'User context not available'}) 

1973 

1974 try: 

1975 user_id = __user__.get('id') 

1976 

1977 # Get the parent message 

1978 parent_message = await Messages.get_message_by_id(parent_message_id) 

1979 

1980 if not parent_message: 

1981 return JSONCodec.dumps({'error': 'Message not found'}) 

1982 

1983 # Verify user has access to the channel 

1984 channel = await Channels.get_channel_by_id(parent_message.channel_id) 

1985 if not channel: 

1986 return JSONCodec.dumps({'error': 'Channel not found'}) 

1987 

1988 user_channels = await Channels.get_channels_by_user_id(user_id) 

1989 channel_ids = [c.id for c in user_channels] 

1990 

1991 if parent_message.channel_id not in channel_ids: 

1992 return JSONCodec.dumps({'error': 'Access denied'}) 

1993 

1994 # Get all thread replies 

1995 thread_replies = await Messages.get_thread_replies_by_message_id(parent_message_id) 

1996 

1997 # Build the response 

1998 messages = [] 

1999 

2000 # Add parent message first 

2001 messages.append( 

2002 { 

2003 'id': parent_message.id, 

2004 'content': parent_message.content, 

2005 'user_id': parent_message.user_id, 

2006 'user_name': parent_message.user.name if parent_message.user else None, 

2007 'is_parent': True, 

2008 'created_at': parent_message.created_at, 

2009 } 

2010 ) 

2011 

2012 # Add thread replies (reverse to get chronological order) 

2013 for reply in reversed(thread_replies): 

2014 messages.append( 

2015 { 

2016 'id': reply.id, 

2017 'content': reply.content, 

2018 'user_id': reply.user_id, 

2019 'user_name': reply.user.name if reply.user else None, 

2020 'is_parent': False, 

2021 'reply_to_id': reply.reply_to_id, 

2022 'created_at': reply.created_at, 

2023 } 

2024 ) 

2025 

2026 return JSONCodec.dumps( 

2027 { 

2028 'channel_id': parent_message.channel_id, 

2029 'channel_name': channel.name, 

2030 'thread_id': parent_message_id, 

2031 'message_count': len(messages), 

2032 'messages': messages, 

2033 }, 

2034 ensure_ascii=False, 

2035 ) 

2036 except Exception as e: 

2037 log.exception(f'view_channel_thread error: {e}') 

2038 return JSONCodec.dumps({'error': str(e)}) 

2039 

2040 

2041# ============================================================================= 

2042# KNOWLEDGE BASE TOOLS 

2043# ============================================================================= 

2044 

2045 

2046async def list_knowledge_bases( 

2047 count: int = 10, 

2048 skip: int = 0, 

2049 __request__: Request = None, 

2050 __user__: dict = None, 

2051) -> str: 

2052 """ 

2053 List the user's accessible knowledge bases so a relevant internal source 

2054 can be chosen. 

2055 

2056 :param count: Maximum number of KBs to return (default: 10) 

2057 :param skip: Number of results to skip for pagination (default: 0) 

2058 :return: JSON with KBs containing id, name, description, and file_count 

2059 """ 

2060 if __request__ is None: 

2061 return JSONCodec.dumps({'error': 'Request context not available'}) 

2062 

2063 if not __user__: 

2064 return JSONCodec.dumps({'error': 'User context not available'}) 

2065 

2066 try: 

2067 from open_webui.models.knowledge import Knowledges 

2068 

2069 user_id = __user__.get('id') 

2070 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

2071 

2072 result = await Knowledges.search_knowledge_bases( 

2073 user_id, 

2074 filter={ 

2075 'query': '', 

2076 'user_id': user_id, 

2077 'group_ids': user_group_ids, 

2078 }, 

2079 skip=skip, 

2080 limit=count, 

2081 ) 

2082 

2083 knowledge_bases = [] 

2084 for knowledge_base in result.items: 

2085 files = await Knowledges.get_files_by_id(knowledge_base.id) 

2086 file_count = len(files) if files else 0 

2087 

2088 knowledge_bases.append( 

2089 { 

2090 'id': knowledge_base.id, 

2091 'name': knowledge_base.name, 

2092 'description': knowledge_base.description or '', 

2093 'file_count': file_count, 

2094 'updated_at': knowledge_base.updated_at, 

2095 } 

2096 ) 

2097 

2098 return JSONCodec.dumps(knowledge_bases, ensure_ascii=False) 

2099 except Exception as e: 

2100 log.exception(f'list_knowledge_bases error: {e}') 

2101 return JSONCodec.dumps({'error': str(e)}) 

2102 

2103 

2104async def search_knowledge_bases( 

2105 query: str, 

2106 count: int = 5, 

2107 skip: int = 0, 

2108 __request__: Request = None, 

2109 __user__: dict = None, 

2110) -> str: 

2111 """ 

2112 Search the user's accessible knowledge bases by name and description to find 

2113 a relevant internal source. 

2114 

2115 :param query: The search query to find matching knowledge bases 

2116 :param count: Maximum number of results to return (default: 5) 

2117 :param skip: Number of results to skip for pagination (default: 0) 

2118 :return: JSON with matching KBs containing id, name, description, and file_count 

2119 """ 

2120 if __request__ is None: 

2121 return JSONCodec.dumps({'error': 'Request context not available'}) 

2122 

2123 if not __user__: 

2124 return JSONCodec.dumps({'error': 'User context not available'}) 

2125 

2126 try: 

2127 from open_webui.models.knowledge import Knowledges 

2128 

2129 user_id = __user__.get('id') 

2130 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

2131 

2132 result = await Knowledges.search_knowledge_bases( 

2133 user_id, 

2134 filter={ 

2135 'query': query, 

2136 'user_id': user_id, 

2137 'group_ids': user_group_ids, 

2138 }, 

2139 skip=skip, 

2140 limit=count, 

2141 ) 

2142 

2143 knowledge_bases = [] 

2144 for knowledge_base in result.items: 

2145 files = await Knowledges.get_files_by_id(knowledge_base.id) 

2146 file_count = len(files) if files else 0 

2147 

2148 knowledge_bases.append( 

2149 { 

2150 'id': knowledge_base.id, 

2151 'name': knowledge_base.name, 

2152 'description': knowledge_base.description or '', 

2153 'file_count': file_count, 

2154 'updated_at': knowledge_base.updated_at, 

2155 } 

2156 ) 

2157 

2158 return JSONCodec.dumps(knowledge_bases, ensure_ascii=False) 

2159 except Exception as e: 

2160 log.exception(f'search_knowledge_bases error: {e}') 

2161 return JSONCodec.dumps({'error': str(e)}) 

2162 

2163 

2164async def search_knowledge_files( 

2165 query: str, 

2166 knowledge_id: Optional[str] = None, 

2167 count: int = 5, 

2168 skip: int = 0, 

2169 __request__: Request = None, 

2170 __user__: dict = None, 

2171 __model_knowledge__: Optional[list[dict]] = None, 

2172) -> str: 

2173 """ 

2174 Search files by filename across knowledge bases the user has access to. 

2175 When the model has attached knowledge, searches only within attached KBs and files. 

2176 Helpful when looking for a specific document or file name. 

2177 

2178 :param query: The search query to find matching files by filename 

2179 :param knowledge_id: Optional KB id to limit search to a specific knowledge base 

2180 :param count: Maximum number of results to return (default: 5) 

2181 :param skip: Number of results to skip for pagination (default: 0) 

2182 :return: JSON with matching files containing id, filename, and updated_at 

2183 """ 

2184 if __request__ is None: 

2185 return JSONCodec.dumps({'error': 'Request context not available'}) 

2186 

2187 if not __user__: 

2188 return JSONCodec.dumps({'error': 'User context not available'}) 

2189 

2190 try: 

2191 from open_webui.models.access_grants import AccessGrants 

2192 from open_webui.models.files import Files 

2193 from open_webui.models.knowledge import Knowledges 

2194 

2195 user_id = __user__.get('id') 

2196 user_role = __user__.get('role', 'user') 

2197 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

2198 

2199 # When model has attached knowledge, scope to attached KBs/files only 

2200 if __model_knowledge__: 

2201 attached_kb_ids = set() 

2202 attached_file_ids = set() 

2203 

2204 for item in __model_knowledge__: 

2205 item_type = item.get('type') 

2206 item_id = item.get('id') 

2207 if item_type == 'collection': 

2208 attached_kb_ids.add(item_id) 

2209 elif item_type == 'file': 

2210 attached_file_ids.add(item_id) 

2211 

2212 # If knowledge_id specified, verify it's in the attached set 

2213 if knowledge_id: 

2214 if knowledge_id not in attached_kb_ids: 

2215 return JSONCodec.dumps({'error': f'Knowledge base {knowledge_id} is not attached to this model'}) 

2216 attached_kb_ids = {knowledge_id} 

2217 

2218 all_files = [] 

2219 

2220 # Search within attached KBs 

2221 for kb_id in attached_kb_ids: 

2222 knowledge = await Knowledges.get_knowledge_by_id(kb_id) 

2223 if not knowledge: 

2224 continue 

2225 

2226 if not ( 

2227 user_role == 'admin' 

2228 or knowledge.user_id == user_id 

2229 or await AccessGrants.has_access( 

2230 user_id=user_id, 

2231 resource_type='knowledge', 

2232 resource_id=knowledge.id, 

2233 permission='read', 

2234 user_group_ids=set(user_group_ids), 

2235 ) 

2236 ): 

2237 continue 

2238 

2239 result = await Knowledges.search_files_by_id( 

2240 knowledge_id=kb_id, 

2241 user_id=user_id, 

2242 filter={'query': query}, 

2243 skip=0, 

2244 limit=count + skip, 

2245 ) 

2246 

2247 for file in result.items: 

2248 all_files.append( 

2249 { 

2250 'id': file.id, 

2251 'filename': file.filename, 

2252 'knowledge_id': knowledge.id, 

2253 'knowledge_name': knowledge.name, 

2254 'updated_at': file.updated_at, 

2255 } 

2256 ) 

2257 

2258 # Search within directly attached files (filename match) 

2259 if not knowledge_id and attached_file_ids: 

2260 query_lower = query.lower() if query else '' 

2261 for file_id in attached_file_ids: 

2262 file = await Files.get_file_by_id(file_id) 

2263 if file and (not query_lower or query_lower in file.filename.lower()): 

2264 all_files.append( 

2265 { 

2266 'id': file.id, 

2267 'filename': file.filename, 

2268 'updated_at': file.updated_at, 

2269 } 

2270 ) 

2271 

2272 # Apply pagination across combined results 

2273 all_files = all_files[skip : skip + count] 

2274 return JSONCodec.dumps(all_files, ensure_ascii=False) 

2275 

2276 # No attached knowledge - search all accessible KBs 

2277 if knowledge_id: 

2278 # search_files_by_id does not enforce knowledge_id ownership; mirror the attached-KB check above. 

2279 knowledge = await Knowledges.get_knowledge_by_id(knowledge_id) 

2280 if not knowledge or not ( 

2281 user_role == 'admin' 

2282 or knowledge.user_id == user_id 

2283 or await AccessGrants.has_access( 

2284 user_id=user_id, 

2285 resource_type='knowledge', 

2286 resource_id=knowledge.id, 

2287 permission='read', 

2288 user_group_ids=set(user_group_ids), 

2289 ) 

2290 ): 

2291 return JSONCodec.dumps({'error': f'Access denied to knowledge base {knowledge_id}'}) 

2292 

2293 result = await Knowledges.search_files_by_id( 

2294 knowledge_id=knowledge_id, 

2295 user_id=user_id, 

2296 filter={'query': query}, 

2297 skip=skip, 

2298 limit=count, 

2299 ) 

2300 else: 

2301 result = await Knowledges.search_knowledge_files( 

2302 filter={ 

2303 'query': query, 

2304 'user_id': user_id, 

2305 'group_ids': user_group_ids, 

2306 }, 

2307 skip=skip, 

2308 limit=count, 

2309 ) 

2310 

2311 files = [] 

2312 for file in result.items: 

2313 file_info = { 

2314 'id': file.id, 

2315 'filename': file.filename, 

2316 'updated_at': file.updated_at, 

2317 } 

2318 if hasattr(file, 'collection') and file.collection: 

2319 file_info['knowledge_id'] = file.collection.get('id', '') 

2320 file_info['knowledge_name'] = file.collection.get('name', '') 

2321 files.append(file_info) 

2322 

2323 return JSONCodec.dumps(files, ensure_ascii=False) 

2324 except Exception as e: 

2325 log.exception(f'search_knowledge_files error: {e}') 

2326 return JSONCodec.dumps({'error': str(e)}) 

2327 

2328 

2329async def _get_accessible_chat_files( 

2330 files: Optional[list[dict]], 

2331 user: dict, 

2332 file_id: Optional[str] = None, 

2333) -> list[tuple[dict, object]]: 

2334 from open_webui.models.files import Files 

2335 

2336 accessible = [] 

2337 seen = set() 

2338 

2339 for item in files or []: 

2340 if not isinstance(item, dict) or item.get('type', 'file') != 'file': 

2341 continue 

2342 fid = item.get('id') or item.get('url') or '' 

2343 if ( 

2344 not isinstance(fid, str) 

2345 or not fid 

2346 or fid in seen 

2347 or fid.startswith(('http://', 'https://', 'data:')) 

2348 or (file_id and fid != file_id) 

2349 ): 

2350 continue 

2351 normalized = {**item, 'id': fid, 'type': 'file'} 

2352 if 'name' not in normalized and item.get('filename'): 

2353 normalized['name'] = item.get('filename') 

2354 seen.add(fid) 

2355 

2356 file = await Files.get_file_by_id(fid) 

2357 if file and await _has_read_access_to_file(file, user): 

2358 accessible.append((normalized, file)) 

2359 

2360 return accessible 

2361 

2362 

2363def _grep_file_models( 

2364 files_to_search: list, 

2365 pattern: str, 

2366 case_insensitive: bool = False, 

2367 count_only: bool = False, 

2368) -> str: 

2369 from open_webui.tools.knowledge_fs import build_matcher 

2370 

2371 matches, err = build_matcher(pattern, case_insensitive) 

2372 if err: 

2373 return JSONCodec.dumps({'error': err}) 

2374 

2375 results = [] 

2376 total_matches = 0 

2377 counts = [] 

2378 

2379 for file in files_to_search: 

2380 content = '' 

2381 if file.data: 

2382 content = file.data.get('content', '') 

2383 if not content: 

2384 continue 

2385 

2386 lines = content.split('\n') 

2387 file_matches = 0 

2388 

2389 for i, line in enumerate(lines, 1): 

2390 if matches(line): 

2391 file_matches += 1 

2392 total_matches += 1 

2393 if not count_only and len(results) < KNOWLEDGE_GREP_MAX_MATCHES: 

2394 results.append(f'{file.id} {file.filename}:{i}: {line}') 

2395 

2396 if file_matches > 0 and count_only: 

2397 counts.append(f'{file.id} {file.filename}: {file_matches}') 

2398 

2399 if count_only: 

2400 if not counts: 

2401 return f'No matches for "{pattern}"' 

2402 return '\n'.join(counts) + f'\n[{total_matches} total matches]' 

2403 

2404 if not results: 

2405 return f'No matches for "{pattern}"' 

2406 

2407 output = '\n'.join(results) 

2408 if total_matches > KNOWLEDGE_GREP_MAX_MATCHES: 

2409 output += f'\n[{KNOWLEDGE_GREP_MAX_MATCHES} of {total_matches} matches shown — use file_id to narrow]' 

2410 return output 

2411 

2412 

2413async def list_chat_files( 

2414 __request__: Request = None, 

2415 __user__: dict = None, 

2416 __files__: list[dict] = None, 

2417) -> str: 

2418 """ 

2419 List files attached to the current chat. 

2420 

2421 :return: JSON with attached chat files containing id, filename, content type, size, and updated time when available 

2422 """ 

2423 if __request__ is None: 

2424 return JSONCodec.dumps({'error': 'Request context not available'}) 

2425 

2426 if not __user__: 

2427 return JSONCodec.dumps({'error': 'User context not available'}) 

2428 

2429 try: 

2430 files = [] 

2431 for item, file in await _get_accessible_chat_files(__files__, __user__): 

2432 file_info = { 

2433 'id': file.id, 

2434 'filename': file.filename, 

2435 'name': item.get('name') or file.filename, 

2436 'type': item.get('type', 'file'), 

2437 'updated_at': file.updated_at, 

2438 } 

2439 content_type = item.get('content_type') or (file.meta or {}).get('content_type') 

2440 size = item.get('size') or (file.meta or {}).get('size') 

2441 if content_type: 

2442 file_info['content_type'] = content_type 

2443 if size: 

2444 file_info['size'] = size 

2445 files.append(file_info) 

2446 

2447 return JSONCodec.dumps(files, ensure_ascii=False) 

2448 except Exception as e: 

2449 log.exception(f'list_chat_files error: {e}') 

2450 return JSONCodec.dumps({'error': str(e)}) 

2451 

2452 

2453async def grep_chat_files( 

2454 pattern: str, 

2455 file_id: Optional[str] = None, 

2456 case_insensitive: bool = False, 

2457 count_only: bool = False, 

2458 __request__: Request = None, 

2459 __user__: dict = None, 

2460 __files__: list[dict] = None, 

2461) -> str: 

2462 """ 

2463 Search exact text across files attached to the current chat. 

2464 Pass file_id from the attached_files block to search one file. 

2465 Auto-detected regex uses RE2 syntax; no lookarounds/backreferences, and shorthand classes are ASCII-only. 

2466 

2467 :param pattern: The text pattern to search for 

2468 :param file_id: Optional attached file ID to search within a single file 

2469 :param case_insensitive: If true, ignore case when matching 

2470 :param count_only: If true, return only match counts per file 

2471 :return: Matching lines with file IDs, filenames, and line numbers 

2472 """ 

2473 if __request__ is None: 

2474 return JSONCodec.dumps({'error': 'Request context not available'}) 

2475 

2476 if not __user__: 

2477 return JSONCodec.dumps({'error': 'User context not available'}) 

2478 

2479 if not pattern or not pattern.strip(): 

2480 return JSONCodec.dumps({'error': 'Pattern is required'}) 

2481 

2482 if isinstance(file_id, str) and file_id.lower() in ('none', 'null', ''): 

2483 file_id = None 

2484 

2485 try: 

2486 attached_ids = set() 

2487 for item in __files__ or []: 

2488 if not isinstance(item, dict) or item.get('type', 'file') != 'file': 

2489 continue 

2490 fid = item.get('id') or item.get('url') 

2491 if isinstance(fid, str) and fid and not fid.startswith(('http://', 'https://', 'data:')): 

2492 attached_ids.add(fid) 

2493 

2494 if not attached_ids: 

2495 return JSONCodec.dumps({'error': 'No files are attached to this chat'}) 

2496 if file_id and file_id not in attached_ids: 

2497 return JSONCodec.dumps({'error': 'File not found'}) 

2498 

2499 files_to_search = [file for _, file in await _get_accessible_chat_files(__files__, __user__, file_id)] 

2500 if not files_to_search: 

2501 return JSONCodec.dumps({'error': 'No accessible files found'}) 

2502 

2503 return await asyncio.to_thread(_grep_file_models, files_to_search, pattern, case_insensitive, count_only) 

2504 except Exception as e: 

2505 log.exception(f'grep_chat_files error: {e}') 

2506 return JSONCodec.dumps({'error': str(e)}) 

2507 

2508 

2509async def query_chat_files( 

2510 query: str, 

2511 file_id: Optional[str] = None, 

2512 count: Optional[int] = None, 

2513 __request__: Request = None, 

2514 __user__: dict = None, 

2515 __files__: list[dict] = None, 

2516) -> str: 

2517 """ 

2518 Search files attached to the current chat using semantic/vector search. 

2519 Pass file_id from the attached_files block to search one file, or omit it to search all attached chat files. 

2520 

2521 :param query: The search query to find semantically relevant content 

2522 :param file_id: Optional attached file ID to search within a single file 

2523 :param count: Maximum number of results to return, capped by the server RAG top k 

2524 :return: JSON with relevant chunks containing content, source filename, and relevance score 

2525 """ 

2526 if __request__ is None: 

2527 return JSONCodec.dumps({'error': 'Request context not available'}) 

2528 

2529 if not __user__: 

2530 return JSONCodec.dumps({'error': 'User context not available'}) 

2531 

2532 if isinstance(file_id, str) and file_id.lower() in ('none', 'null', ''): 

2533 file_id = None 

2534 if isinstance(count, str): 

2535 if count.lower() in ('none', 'null', ''): 

2536 count = None 

2537 else: 

2538 try: 

2539 count = int(count) 

2540 except ValueError: 

2541 count = None 

2542 

2543 try: 

2544 from open_webui.retrieval.utils import get_sources_from_items 

2545 

2546 attached_ids = set() 

2547 for item in __files__ or []: 

2548 if not isinstance(item, dict) or item.get('type', 'file') != 'file': 

2549 continue 

2550 fid = item.get('id') or item.get('url') 

2551 if isinstance(fid, str) and fid and not fid.startswith(('http://', 'https://', 'data:')): 

2552 attached_ids.add(fid) 

2553 

2554 if not attached_ids: 

2555 return JSONCodec.dumps({'error': 'No files are attached to this chat'}) 

2556 if file_id and file_id not in attached_ids: 

2557 return JSONCodec.dumps({'error': 'File not found'}) 

2558 

2559 accessible = await _get_accessible_chat_files(__files__, __user__, file_id) 

2560 if not accessible: 

2561 return JSONCodec.dumps({'error': 'No accessible files found'}) 

2562 

2563 file_items = [{**item} for item, _ in accessible] 

2564 rag_config = await Config.get_many( 

2565 'rag.top_k', 

2566 'rag.top_k_reranker', 

2567 'rag.relevance_threshold', 

2568 'rag.hybrid_bm25_weight', 

2569 'rag.enable_hybrid_search', 

2570 'rag.full_context', 

2571 ) 

2572 top_k = rag_config.get('rag.top_k') or 5 

2573 count = top_k if count is None else max(1, min(count, top_k)) 

2574 full_context = all(item.get('context') == 'full' for item in file_items) or rag_config.get('rag.full_context') 

2575 

2576 embedding_function = getattr(__request__.app.state, 'EMBEDDING_FUNCTION', None) 

2577 if not embedding_function and not full_context: 

2578 return JSONCodec.dumps({'error': 'Embedding function not configured'}) 

2579 

2580 user_model = UserModel(**__user__) 

2581 sources = await get_sources_from_items( 

2582 request=__request__, 

2583 items=file_items, 

2584 queries=[query], 

2585 embedding_function=( 

2586 lambda queries, prefix: ( 

2587 embedding_function(queries, prefix=prefix, user=user_model) if embedding_function else None 

2588 ) 

2589 ), 

2590 k=count, 

2591 reranking_function=( 

2592 (lambda q, docs: __request__.app.state.RERANKING_FUNCTION(q, docs, user=user_model)) 

2593 if getattr(__request__.app.state, 'RERANKING_FUNCTION', None) 

2594 else None 

2595 ), 

2596 k_reranker=rag_config.get('rag.top_k_reranker'), 

2597 r=rag_config.get('rag.relevance_threshold'), 

2598 hybrid_bm25_weight=rag_config.get('rag.hybrid_bm25_weight'), 

2599 hybrid_search=rag_config.get('rag.enable_hybrid_search'), 

2600 full_context=full_context, 

2601 user=user_model, 

2602 ) 

2603 

2604 chunks = [] 

2605 for source in sources or []: 

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

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

2608 distances = source.get('distances') or [] 

2609 source_info = source.get('source') or {} 

2610 

2611 for idx, doc in enumerate(documents): 

2612 metadata = metadatas[idx] if idx < len(metadatas) and isinstance(metadatas[idx], dict) else {} 

2613 chunk = { 

2614 **filter_source_metadata(metadata), 

2615 'content': doc, 

2616 'source': metadata.get('source', metadata.get('name', source_info.get('name', 'Unknown'))), 

2617 'file_id': metadata.get('file_id', source_info.get('id', '')), 

2618 } 

2619 if idx < len(distances): 

2620 chunk['distance'] = distances[idx] 

2621 chunks.append(chunk) 

2622 

2623 return JSONCodec.dumps(chunks[:count], ensure_ascii=False) 

2624 except Exception as e: 

2625 log.exception(f'query_chat_files error: {e}') 

2626 return JSONCodec.dumps({'error': str(e)}) 

2627 

2628 

2629async def grep_knowledge_files( 

2630 pattern: str, 

2631 file_id: Optional[str] = None, 

2632 case_insensitive: bool = False, 

2633 count_only: bool = False, 

2634 __request__: Request = None, 

2635 __user__: dict = None, 

2636 __model_knowledge__: Optional[list[dict]] = None, 

2637) -> str: 

2638 """ 

2639 Search for exact text across knowledge files. Returns matching lines with line numbers. 

2640 Unlike query_knowledge_files (semantic/vector search), this performs exact string matching. 

2641 Automatically detects regex patterns (e.g. "error|warn", "version \\d+"). 

2642 Regex uses RE2 syntax; no lookarounds/backreferences, and shorthand character classes are ASCII-only. 

2643 Helpful for literal strings, identifiers, error messages, or regex-style searches. 

2644 

2645 :param pattern: The text pattern to search for (regex auto-detected) 

2646 :param file_id: Optional file ID to search within a single file only 

2647 :param case_insensitive: If true, ignore case when matching (default: false) 

2648 :param count_only: If true, return only match counts per file (default: false) 

2649 :return: Matching lines with file IDs, filenames, and line numbers 

2650 """ 

2651 if __request__ is None: 

2652 return JSONCodec.dumps({'error': 'Request context not available'}) 

2653 

2654 if not __user__: 

2655 return JSONCodec.dumps({'error': 'User context not available'}) 

2656 

2657 if not pattern or not pattern.strip(): 

2658 return JSONCodec.dumps({'error': 'Pattern is required'}) 

2659 

2660 try: 

2661 from open_webui.models.files import Files 

2662 from open_webui.models.knowledge import Knowledges 

2663 

2664 user_id = __user__.get('id') 

2665 user_role = __user__.get('role', 'user') 

2666 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

2667 

2668 # Collect files to search 

2669 files_to_search = [] 

2670 

2671 if file_id: 

2672 # Single file mode — verify access 

2673 file = await Files.get_file_by_id(file_id) 

2674 if file: 

2675 if not await _has_read_access_to_file(file, __user__, __model_knowledge__): 

2676 return JSONCodec.dumps({'error': 'File not found'}) 

2677 files_to_search.append(file) 

2678 elif __model_knowledge__: 

2679 # Scoped to model's attached knowledge 

2680 from open_webui.models.access_grants import AccessGrants 

2681 

2682 seen_ids = set() 

2683 for item in __model_knowledge__: 

2684 item_type = item.get('type') 

2685 item_id = item.get('id') 

2686 if item_type == 'file' and item_id not in seen_ids: 

2687 file = await Files.get_file_by_id(item_id) 

2688 if file: 

2689 files_to_search.append(file) 

2690 seen_ids.add(item_id) 

2691 elif item_type == 'collection': 

2692 knowledge = await Knowledges.get_knowledge_by_id(item_id) 

2693 if not knowledge: 

2694 continue 

2695 # Verify user can access this KB 

2696 if not ( 

2697 user_role == 'admin' 

2698 or knowledge.user_id == user_id 

2699 or await AccessGrants.has_access( 

2700 user_id=user_id, 

2701 resource_type='knowledge', 

2702 resource_id=knowledge.id, 

2703 permission='read', 

2704 user_group_ids=set(user_group_ids), 

2705 ) 

2706 ): 

2707 continue 

2708 kb_files = await Knowledges.get_files_by_id(item_id) 

2709 if kb_files: 

2710 for f in kb_files: 

2711 if f.id not in seen_ids: 

2712 files_to_search.append(f) 

2713 seen_ids.add(f.id) 

2714 else: 

2715 # All accessible knowledge bases — use the same search pattern as list_knowledge_bases 

2716 result = await Knowledges.search_knowledge_bases( 

2717 user_id, 

2718 filter={ 

2719 'query': '', 

2720 'user_id': user_id, 

2721 'group_ids': user_group_ids, 

2722 }, 

2723 skip=0, 

2724 limit=200, 

2725 ) 

2726 seen_ids = set() 

2727 for kb in result.items: 

2728 file_ids = [] 

2729 # Get files attached to this KB 

2730 files_from_kb = await Knowledges.get_files_by_id(kb.id) 

2731 if files_from_kb: 

2732 file_ids = [f.id for f in files_from_kb] 

2733 for fid in file_ids: 

2734 if fid not in seen_ids: 

2735 file = await Files.get_file_by_id(fid) 

2736 if file: 

2737 files_to_search.append(file) 

2738 seen_ids.add(fid) 

2739 

2740 if not files_to_search: 

2741 return JSONCodec.dumps({'error': 'No accessible files found'}) 

2742 

2743 return await asyncio.to_thread(_grep_file_models, files_to_search, pattern, case_insensitive, count_only) 

2744 

2745 except Exception as e: 

2746 log.exception(f'grep_knowledge_files error: {e}') 

2747 return JSONCodec.dumps({'error': str(e)}) 

2748 

2749 

2750async def view_file( 

2751 file_id: str, 

2752 offset: int = 0, 

2753 max_chars: int = VIEW_FILE_DEFAULT_MAX_CHARS, 

2754 line_numbers: bool = False, 

2755 start_line: Optional[int] = None, 

2756 end_line: Optional[int] = None, 

2757 __request__: Request = None, 

2758 __user__: dict = None, 

2759 __model_knowledge__: Optional[list[dict]] = None, 

2760) -> str: 

2761 """ 

2762 Get the content of a file by its ID. Supports pagination for large files. 

2763 

2764 :param file_id: The ID of the file to retrieve 

2765 :param offset: Character offset to start reading from (default: 0) 

2766 :param max_chars: Maximum characters to return (a server-side hard cap applies) 

2767 :param line_numbers: If true, prefix each line with its 1-indexed line number 

2768 :param start_line: Optional 1-indexed start line (overrides offset/max_chars when set) 

2769 :param end_line: Optional 1-indexed end line (inclusive) 

2770 :return: JSON with the file's id, filename, content, and pagination metadata if truncated 

2771 """ 

2772 if __request__ is None: 

2773 return JSONCodec.dumps({'error': 'Request context not available'}) 

2774 

2775 if not __user__: 

2776 return JSONCodec.dumps({'error': 'User context not available'}) 

2777 

2778 # Coerce parameters from LLM tool calls (may come as strings) 

2779 if isinstance(offset, str): 

2780 try: 

2781 offset = int(offset) 

2782 except ValueError: 

2783 offset = 0 

2784 if isinstance(max_chars, str): 

2785 try: 

2786 max_chars = int(max_chars) 

2787 except ValueError: 

2788 max_chars = VIEW_FILE_DEFAULT_MAX_CHARS 

2789 

2790 # Enforce hard cap 

2791 max_chars = min(max(max_chars, 1), VIEW_FILE_MAX_CHARS) 

2792 offset = max(offset, 0) 

2793 

2794 try: 

2795 from open_webui.models.files import Files 

2796 

2797 file = await Files.get_file_by_id(file_id) 

2798 if not file: 

2799 return JSONCodec.dumps({'error': 'File not found'}) 

2800 

2801 if not await _has_read_access_to_file(file, __user__, __model_knowledge__): 

2802 return JSONCodec.dumps({'error': 'File not found'}) 

2803 

2804 content = '' 

2805 if file.data: 

2806 content = file.data.get('content', '') 

2807 

2808 total_chars = len(content) 

2809 

2810 # Line-based addressing (overrides char-based offset/max_chars) 

2811 if start_line is not None: 

2812 all_lines = content.split('\n') 

2813 total_lines = len(all_lines) 

2814 s = max(1, int(start_line)) - 1 # 1-indexed to 0-indexed 

2815 e = min(total_lines, int(end_line) if end_line else s + 100) 

2816 selected = all_lines[s:e] 

2817 sliced = '\n'.join(f'{s + i + 1}: {line}' for i, line in enumerate(selected)) 

2818 is_truncated = e < total_lines 

2819 result = { 

2820 'id': file.id, 

2821 'filename': file.filename, 

2822 'content': sliced, 

2823 'updated_at': file.updated_at, 

2824 'created_at': file.created_at, 

2825 'total_lines': total_lines, 

2826 'showing_lines': f'{s + 1}-{e}', 

2827 } 

2828 if is_truncated: 

2829 result['truncated'] = True 

2830 result['next_start_line'] = e + 1 

2831 return JSONCodec.dumps(result, ensure_ascii=False) 

2832 

2833 sliced = content[offset : offset + max_chars] 

2834 is_truncated = (offset + len(sliced)) < total_chars 

2835 

2836 if line_numbers: 

2837 start_ln = content[:offset].count('\n') + 1 

2838 lines = sliced.split('\n') 

2839 sliced = '\n'.join(f'{start_ln + i}: {line}' for i, line in enumerate(lines)) 

2840 

2841 result = { 

2842 'id': file.id, 

2843 'filename': file.filename, 

2844 'content': sliced, 

2845 'updated_at': file.updated_at, 

2846 'created_at': file.created_at, 

2847 } 

2848 

2849 if is_truncated or offset > 0: 

2850 result['truncated'] = is_truncated 

2851 result['total_chars'] = total_chars 

2852 result['returned_chars'] = len(sliced) 

2853 result['offset'] = offset 

2854 if is_truncated: 

2855 result['next_offset'] = offset + len(sliced) 

2856 

2857 return JSONCodec.dumps(result, ensure_ascii=False) 

2858 except Exception as e: 

2859 log.exception(f'view_file error: {e}') 

2860 return JSONCodec.dumps({'error': str(e)}) 

2861 

2862 

2863async def view_knowledge_file( 

2864 file_id: str, 

2865 offset: int = 0, 

2866 max_chars: int = VIEW_FILE_DEFAULT_MAX_CHARS, 

2867 line_numbers: bool = False, 

2868 start_line: Optional[int] = None, 

2869 end_line: Optional[int] = None, 

2870 __request__: Request = None, 

2871 __user__: dict = None, 

2872) -> str: 

2873 """ 

2874 Get the content of a file from a knowledge base. Supports pagination for large files. 

2875 

2876 :param file_id: The ID of the file to retrieve 

2877 :param offset: Character offset to start reading from (default: 0) 

2878 :param max_chars: Maximum characters to return (a server-side hard cap applies) 

2879 :param line_numbers: If true, prefix each line with its 1-indexed line number 

2880 :param start_line: Optional 1-indexed start line (overrides offset/max_chars when set) 

2881 :param end_line: Optional 1-indexed end line (inclusive) 

2882 :return: JSON with the file's id, filename, content, and pagination metadata if truncated 

2883 """ 

2884 if __request__ is None: 

2885 return JSONCodec.dumps({'error': 'Request context not available'}) 

2886 

2887 if not __user__: 

2888 return JSONCodec.dumps({'error': 'User context not available'}) 

2889 

2890 # Coerce parameters from LLM tool calls (may come as strings) 

2891 if isinstance(offset, str): 

2892 try: 

2893 offset = int(offset) 

2894 except ValueError: 

2895 offset = 0 

2896 if isinstance(max_chars, str): 

2897 try: 

2898 max_chars = int(max_chars) 

2899 except ValueError: 

2900 max_chars = VIEW_FILE_DEFAULT_MAX_CHARS 

2901 

2902 # Enforce hard cap 

2903 max_chars = min(max(max_chars, 1), VIEW_FILE_MAX_CHARS) 

2904 offset = max(offset, 0) 

2905 

2906 try: 

2907 from open_webui.models.access_grants import AccessGrants 

2908 from open_webui.models.files import Files 

2909 from open_webui.models.knowledge import Knowledges 

2910 

2911 user_id = __user__.get('id') 

2912 user_role = __user__.get('role', 'user') 

2913 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

2914 

2915 file = await Files.get_file_by_id(file_id) 

2916 if not file: 

2917 return JSONCodec.dumps({'error': 'File not found'}) 

2918 

2919 # Check access via any KB containing this file 

2920 knowledges = await Knowledges.get_knowledges_by_file_id(file_id) 

2921 has_knowledge_access = False 

2922 knowledge_info = None 

2923 

2924 for knowledge_base in knowledges: 

2925 if ( 

2926 user_role == 'admin' 

2927 or knowledge_base.user_id == user_id 

2928 or await AccessGrants.has_access( 

2929 user_id=user_id, 

2930 resource_type='knowledge', 

2931 resource_id=knowledge_base.id, 

2932 permission='read', 

2933 user_group_ids=set(user_group_ids), 

2934 ) 

2935 ): 

2936 has_knowledge_access = True 

2937 knowledge_info = {'id': knowledge_base.id, 'name': knowledge_base.name} 

2938 break 

2939 

2940 if not has_knowledge_access: 

2941 if file.user_id != user_id and user_role != 'admin': 

2942 return JSONCodec.dumps({'error': 'Access denied'}) 

2943 

2944 content = '' 

2945 if file.data: 

2946 content = file.data.get('content', '') 

2947 

2948 total_chars = len(content) 

2949 

2950 # Line-based addressing (overrides char-based offset/max_chars) 

2951 if start_line is not None: 

2952 all_lines = content.split('\n') 

2953 total_lines = len(all_lines) 

2954 s = max(1, int(start_line)) - 1 

2955 e = min(total_lines, int(end_line) if end_line else s + 100) 

2956 selected = all_lines[s:e] 

2957 sliced = '\n'.join(f'{s + i + 1}: {line}' for i, line in enumerate(selected)) 

2958 is_truncated = e < total_lines 

2959 result = { 

2960 'id': file.id, 

2961 'filename': file.filename, 

2962 'content': sliced, 

2963 'updated_at': file.updated_at, 

2964 'created_at': file.created_at, 

2965 'total_lines': total_lines, 

2966 'showing_lines': f'{s + 1}-{e}', 

2967 } 

2968 if knowledge_info: 

2969 result['knowledge_id'] = knowledge_info['id'] 

2970 result['knowledge_name'] = knowledge_info['name'] 

2971 if is_truncated: 

2972 result['truncated'] = True 

2973 result['next_start_line'] = e + 1 

2974 return JSONCodec.dumps(result, ensure_ascii=False) 

2975 

2976 sliced = content[offset : offset + max_chars] 

2977 is_truncated = (offset + len(sliced)) < total_chars 

2978 

2979 if line_numbers: 

2980 start_ln = content[:offset].count('\n') + 1 

2981 lines = sliced.split('\n') 

2982 sliced = '\n'.join(f'{start_ln + i}: {line}' for i, line in enumerate(lines)) 

2983 

2984 result = { 

2985 'id': file.id, 

2986 'filename': file.filename, 

2987 'content': sliced, 

2988 'updated_at': file.updated_at, 

2989 'created_at': file.created_at, 

2990 } 

2991 if knowledge_info: 

2992 result['knowledge_id'] = knowledge_info['id'] 

2993 result['knowledge_name'] = knowledge_info['name'] 

2994 

2995 if is_truncated or offset > 0: 

2996 result['truncated'] = is_truncated 

2997 result['total_chars'] = total_chars 

2998 result['returned_chars'] = len(sliced) 

2999 result['offset'] = offset 

3000 if is_truncated: 

3001 result['next_offset'] = offset + len(sliced) 

3002 

3003 return JSONCodec.dumps(result, ensure_ascii=False) 

3004 except Exception as e: 

3005 log.exception(f'view_knowledge_file error: {e}') 

3006 return JSONCodec.dumps({'error': str(e)}) 

3007 

3008 

3009async def list_knowledge( 

3010 knowledge_id: Optional[str] = None, 

3011 skip: int = 0, 

3012 count: int = 50, 

3013 __request__: Request = None, 

3014 __user__: dict = None, 

3015 __model_knowledge__: Optional[list[dict]] = None, 

3016) -> str: 

3017 """ 

3018 List knowledge bases, files, and notes attached to the current model. 

3019 Use this first to discover what knowledge is available before querying or reading files. 

3020 Without knowledge_id: returns KB summaries (name, description, file_count) 

3021 plus standalone files and notes — no file listing inside KBs. 

3022 With knowledge_id: includes paginated file listing for that specific KB. 

3023 Use skip/count to page through large KBs. 

3024 

3025 :param knowledge_id: Optional KB ID to get file listing for 

3026 :param skip: Number of files to skip for pagination (default: 0) 

3027 :param count: Maximum files per page (default: 50, max: 200) 

3028 :return: JSON with knowledge_bases, files, and notes attached to this model 

3029 """ 

3030 if __request__ is None: 

3031 return JSONCodec.dumps({'error': 'Request context not available'}) 

3032 

3033 if not __user__: 

3034 return JSONCodec.dumps({'error': 'User context not available'}) 

3035 

3036 if not __model_knowledge__: 

3037 return JSONCodec.dumps({'knowledge_bases': [], 'files': [], 'notes': []}) 

3038 

3039 # Coerce parameters from LLM tool calls (may come as strings) 

3040 if isinstance(skip, str): 

3041 try: 

3042 skip = int(skip) 

3043 except ValueError: 

3044 skip = 0 

3045 if isinstance(count, str): 

3046 try: 

3047 count = int(count) 

3048 except ValueError: 

3049 count = 50 

3050 if isinstance(knowledge_id, str) and knowledge_id.lower() in ('none', 'null', ''): 

3051 knowledge_id = None 

3052 

3053 count = min(count, 200) 

3054 

3055 try: 

3056 from open_webui.models.access_grants import AccessGrants 

3057 from open_webui.models.files import Files 

3058 from open_webui.models.knowledge import Knowledges 

3059 from open_webui.models.notes import Notes 

3060 

3061 user_id = __user__.get('id') 

3062 user_role = __user__.get('role', 'user') 

3063 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

3064 

3065 knowledge_bases = [] 

3066 files = [] 

3067 notes = [] 

3068 

3069 for item in __model_knowledge__: 

3070 item_type = item.get('type') 

3071 item_id = item.get('id') 

3072 

3073 if item_type == 'collection': 

3074 knowledge = await Knowledges.get_knowledge_by_id(item_id) 

3075 if knowledge and ( 

3076 user_role == 'admin' 

3077 or knowledge.user_id == user_id 

3078 or await AccessGrants.has_access( 

3079 user_id=user_id, 

3080 resource_type='knowledge', 

3081 resource_id=knowledge.id, 

3082 permission='read', 

3083 user_group_ids=set(user_group_ids), 

3084 ) 

3085 ): 

3086 kb_files = await Knowledges.get_files_by_id(knowledge.id) 

3087 file_count = len(kb_files) if kb_files else 0 

3088 

3089 kb_entry = { 

3090 'id': knowledge.id, 

3091 'name': knowledge.name, 

3092 'description': knowledge.description or '', 

3093 'file_count': file_count, 

3094 } 

3095 

3096 # Include file listing only when this KB is targeted 

3097 if knowledge_id and knowledge_id == knowledge.id: 

3098 if kb_files: 

3099 paged_files = kb_files[skip : skip + count] 

3100 kb_entry['files'] = [{'id': f.id, 'filename': f.filename} for f in paged_files] 

3101 kb_entry['files_skip'] = skip 

3102 kb_entry['files_count'] = len(paged_files) 

3103 kb_entry['files_total'] = file_count 

3104 kb_entry['has_more'] = skip + count < file_count 

3105 

3106 knowledge_bases.append(kb_entry) 

3107 

3108 elif item_type == 'file': 

3109 file = await Files.get_file_by_id(item_id) 

3110 if file: 

3111 files.append( 

3112 { 

3113 'id': file.id, 

3114 'filename': file.filename, 

3115 'updated_at': file.updated_at, 

3116 } 

3117 ) 

3118 

3119 elif item_type == 'note': 

3120 note = await Notes.get_note_by_id(item_id) 

3121 if note and ( 

3122 user_role == 'admin' 

3123 or note.user_id == user_id 

3124 or await AccessGrants.has_access( 

3125 user_id=user_id, 

3126 resource_type='note', 

3127 resource_id=note.id, 

3128 permission='read', 

3129 ) 

3130 ): 

3131 notes.append( 

3132 { 

3133 'id': note.id, 

3134 'title': note.title, 

3135 } 

3136 ) 

3137 

3138 return JSONCodec.dumps( 

3139 { 

3140 'knowledge_bases': knowledge_bases, 

3141 'files': files, 

3142 'notes': notes, 

3143 }, 

3144 ensure_ascii=False, 

3145 ) 

3146 except Exception as e: 

3147 log.exception(f'list_knowledge error: {e}') 

3148 return JSONCodec.dumps({'error': str(e)}) 

3149 

3150 

3151async def query_knowledge_files( 

3152 query: str, 

3153 knowledge_ids: Optional[list[str]] = None, 

3154 count: int = 5, 

3155 __request__: Request = None, 

3156 __user__: dict = None, 

3157 __model_knowledge__: list[dict] = None, 

3158) -> str: 

3159 """ 

3160 Search knowledge base files using semantic/vector search. Searches across collections (KBs), 

3161 individual files, and notes that the user has access to. 

3162 Helpful for internal documentation, uploaded knowledge, and attached model knowledge. 

3163 

3164 :param query: The search query to find semantically relevant content 

3165 :param knowledge_ids: Optional list of KB ids to limit search to specific knowledge bases 

3166 :param count: Maximum number of results to return (default: 5) 

3167 :return: JSON with relevant chunks containing content, source filename, and relevance score 

3168 """ 

3169 if __request__ is None: 

3170 return JSONCodec.dumps({'error': 'Request context not available'}) 

3171 

3172 if not __user__: 

3173 return JSONCodec.dumps({'error': 'User context not available'}) 

3174 

3175 # Coerce parameters from LLM tool calls (may come as strings) 

3176 if isinstance(count, str): 

3177 try: 

3178 count = int(count) 

3179 except ValueError: 

3180 count = 5 # Default fallback 

3181 

3182 # Handle knowledge_ids being string "None", "null", or empty 

3183 if isinstance(knowledge_ids, str): 

3184 if knowledge_ids.lower() in ('none', 'null', ''): 

3185 knowledge_ids = None 

3186 else: 

3187 # Try to parse as JSON array if it looks like one 

3188 try: 

3189 knowledge_ids = JSONCodec.loads(knowledge_ids) 

3190 except JSONCodec.JSONDecodeError: 

3191 # Treat as single ID 

3192 knowledge_ids = [knowledge_ids] 

3193 

3194 try: 

3195 from open_webui.models.access_grants import AccessGrants 

3196 from open_webui.models.files import Files 

3197 from open_webui.models.knowledge import Knowledges 

3198 from open_webui.models.notes import Notes 

3199 from open_webui.retrieval.external import retrieve_external_knowledge 

3200 from open_webui.retrieval.utils import query_collection 

3201 

3202 user_id = __user__.get('id') 

3203 user_role = __user__.get('role', 'user') 

3204 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

3205 

3206 embedding_function = getattr(__request__.app.state, 'EMBEDDING_FUNCTION', None) 

3207 if not embedding_function: 

3208 return JSONCodec.dumps({'error': 'Embedding function not configured'}) 

3209 user_model = UserModel(**__user__) 

3210 

3211 collection_names = [] 

3212 external_knowledges = [] 

3213 note_results = [] # Notes aren't vectorized, handle separately 

3214 

3215 # If model has attached knowledge, use those 

3216 if __model_knowledge__: 

3217 for item in __model_knowledge__: 

3218 item_type = item.get('type') 

3219 item_id = item.get('id') 

3220 

3221 if item_type == 'collection': 

3222 # Knowledge base - use KB ID as collection name 

3223 knowledge = await Knowledges.get_knowledge_by_id(item_id) 

3224 if knowledge and ( 

3225 user_role == 'admin' 

3226 or knowledge.user_id == user_id 

3227 or await AccessGrants.has_access( 

3228 user_id=user_id, 

3229 resource_type='knowledge', 

3230 resource_id=knowledge.id, 

3231 permission='read', 

3232 user_group_ids=set(user_group_ids), 

3233 ) 

3234 ): 

3235 if (knowledge.meta or {}).get('source') == 'external': 

3236 external_knowledges.append(knowledge) 

3237 else: 

3238 collection_names.append(item_id) 

3239 

3240 elif item_type == 'file': 

3241 # Individual file - use file-{id} as collection name 

3242 file = await Files.get_file_by_id(item_id) 

3243 if file: 

3244 collection_names.append(f'file-{item_id}') 

3245 

3246 elif item_type == 'note': 

3247 # Note - always return full content as context 

3248 note = await Notes.get_note_by_id(item_id) 

3249 if note and ( 

3250 user_role == 'admin' 

3251 or note.user_id == user_id 

3252 or await AccessGrants.has_access( 

3253 user_id=user_id, 

3254 resource_type='note', 

3255 resource_id=note.id, 

3256 permission='read', 

3257 ) 

3258 ): 

3259 content = note.data.get('content', {}).get('md', '') 

3260 note_results.append( 

3261 { 

3262 'content': content, 

3263 'source': note.title, 

3264 'note_id': note.id, 

3265 'type': 'note', 

3266 } 

3267 ) 

3268 

3269 elif knowledge_ids: 

3270 # User specified specific KBs 

3271 for knowledge_id in knowledge_ids: 

3272 knowledge = await Knowledges.get_knowledge_by_id(knowledge_id) 

3273 if knowledge and ( 

3274 user_role == 'admin' 

3275 or knowledge.user_id == user_id 

3276 or await AccessGrants.has_access( 

3277 user_id=user_id, 

3278 resource_type='knowledge', 

3279 resource_id=knowledge.id, 

3280 permission='read', 

3281 user_group_ids=set(user_group_ids), 

3282 ) 

3283 ): 

3284 if (knowledge.meta or {}).get('source') == 'external': 

3285 external_knowledges.append(knowledge) 

3286 else: 

3287 collection_names.append(knowledge_id) 

3288 else: 

3289 # No model knowledge and no specific IDs - search all accessible KBs 

3290 result = await Knowledges.search_knowledge_bases( 

3291 user_id, 

3292 filter={ 

3293 'query': '', 

3294 'user_id': user_id, 

3295 'group_ids': user_group_ids, 

3296 }, 

3297 skip=0, 

3298 limit=50, 

3299 ) 

3300 for knowledge_base in result.items: 

3301 if (knowledge_base.meta or {}).get('source') == 'external': 

3302 external_knowledges.append(knowledge_base) 

3303 else: 

3304 collection_names.append(knowledge_base.id) 

3305 

3306 chunks = [] 

3307 

3308 # Add note results first 

3309 chunks.extend(note_results) 

3310 

3311 # Query vector collections if any 

3312 if collection_names: 

3313 query_results = await query_collection( 

3314 __request__, 

3315 collection_names=collection_names, 

3316 queries=[query], 

3317 embedding_function=lambda queries, prefix: embedding_function(queries, prefix=prefix, user=user_model), 

3318 k=count, 

3319 ) 

3320 

3321 if query_results and 'documents' in query_results: 

3322 documents = query_results.get('documents', [[]])[0] 

3323 metadatas = query_results.get('metadatas', [[]])[0] 

3324 distances = query_results.get('distances', [[]])[0] 

3325 

3326 for idx, doc in enumerate(documents): 

3327 chunk_info = { 

3328 **filter_source_metadata(metadatas[idx]), 

3329 'content': doc, 

3330 'source': metadatas[idx].get('source', metadatas[idx].get('name', 'Unknown')), 

3331 'file_id': metadatas[idx].get('file_id', ''), 

3332 } 

3333 if idx < len(distances): 

3334 chunk_info['distance'] = distances[idx] 

3335 chunks.append(chunk_info) 

3336 

3337 for knowledge in external_knowledges: 

3338 query_results = await retrieve_external_knowledge( 

3339 __request__, 

3340 knowledge, 

3341 queries=[query], 

3342 count=count, 

3343 user=user_model, 

3344 ) 

3345 documents = query_results.get('documents', [[]])[0] 

3346 metadatas = query_results.get('metadatas', [[]])[0] 

3347 distances = query_results.get('distances', [[]])[0] 

3348 

3349 for idx, doc in enumerate(documents): 

3350 metadata = metadatas[idx] if idx < len(metadatas) else {} 

3351 chunk_info = { 

3352 **filter_source_metadata(metadata), 

3353 'content': doc, 

3354 'source': metadata.get('source', metadata.get('name', knowledge.name)), 

3355 'file_id': metadata.get('file_id', f'external-{knowledge.id}'), 

3356 'type': 'external', 

3357 'knowledge_id': knowledge.id, 

3358 } 

3359 if idx < len(distances): 

3360 chunk_info['distance'] = distances[idx] 

3361 chunks.append(chunk_info) 

3362 

3363 # Limit to requested count 

3364 chunks = chunks[:count] 

3365 

3366 return JSONCodec.dumps(chunks, ensure_ascii=False) 

3367 except Exception as e: 

3368 log.exception(f'query_knowledge_files error: {e}') 

3369 return JSONCodec.dumps({'error': str(e)}) 

3370 

3371 

3372async def query_knowledge_bases( 

3373 query: str, 

3374 count: int = 5, 

3375 __request__: Request = None, 

3376 __user__: dict = None, 

3377) -> str: 

3378 """ 

3379 Search knowledge bases by semantic similarity to query. 

3380 Finds KBs whose name/description match the meaning of your query. 

3381 Helpful for discovering which knowledge base to query next. 

3382 

3383 :param query: Natural language query describing what you're looking for 

3384 :param count: Maximum results (default: 5) 

3385 :return: JSON with matching KBs (id, name, description, similarity) 

3386 """ 

3387 if __request__ is None: 

3388 return JSONCodec.dumps({'error': 'Request context not available'}) 

3389 

3390 if not __user__: 

3391 return JSONCodec.dumps({'error': 'User context not available'}) 

3392 

3393 try: 

3394 import heapq 

3395 

3396 from open_webui.models.knowledge import Knowledges 

3397 from open_webui.retrieval.vector.async_client import ASYNC_VECTOR_DB_CLIENT 

3398 from open_webui.routers.knowledge import KNOWLEDGE_BASES_COLLECTION 

3399 

3400 user_id = __user__.get('id') 

3401 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

3402 embedding_function = getattr(__request__.app.state, 'EMBEDDING_FUNCTION', None) 

3403 if not embedding_function: 

3404 return JSONCodec.dumps({'error': 'Embedding function not configured'}) 

3405 user_model = UserModel(**__user__) 

3406 query_embedding = await embedding_function(query, prefix=RAG_EMBEDDING_QUERY_PREFIX, user=user_model) 

3407 

3408 # Min-heap of (distance, knowledge_base_id) - only holds top `count` results 

3409 top_results_heap = [] 

3410 seen_ids = set() 

3411 page_offset = 0 

3412 page_size = 100 

3413 

3414 while True: 

3415 accessible_knowledge_bases = await Knowledges.search_knowledge_bases( 

3416 user_id, 

3417 filter={'user_id': user_id, 'group_ids': user_group_ids}, 

3418 skip=page_offset, 

3419 limit=page_size, 

3420 ) 

3421 

3422 if not accessible_knowledge_bases.items: 

3423 break 

3424 

3425 accessible_ids = [kb.id for kb in accessible_knowledge_bases.items] 

3426 

3427 search_results = await ASYNC_VECTOR_DB_CLIENT.search( 

3428 collection_name=KNOWLEDGE_BASES_COLLECTION, 

3429 vectors=[query_embedding], 

3430 filter={'knowledge_base_id': {'$in': accessible_ids}}, 

3431 limit=count, 

3432 ) 

3433 

3434 if search_results and search_results.ids and search_results.ids[0]: 

3435 result_ids = search_results.ids[0] 

3436 result_distances = search_results.distances[0] if search_results.distances else [0] * len(result_ids) 

3437 

3438 for knowledge_base_id, distance in zip(result_ids, result_distances): 

3439 if knowledge_base_id in seen_ids: 

3440 continue 

3441 seen_ids.add(knowledge_base_id) 

3442 

3443 if len(top_results_heap) < count: 

3444 heapq.heappush(top_results_heap, (distance, knowledge_base_id)) 

3445 elif distance > top_results_heap[0][0]: 

3446 heapq.heapreplace(top_results_heap, (distance, knowledge_base_id)) 

3447 

3448 page_offset += page_size 

3449 if len(accessible_knowledge_bases.items) < page_size: 

3450 break 

3451 if page_offset >= MAX_KNOWLEDGE_BASE_SEARCH_ITEMS: 

3452 break 

3453 

3454 # Sort by distance descending (best first) and fetch KB details 

3455 sorted_results = sorted(top_results_heap, key=lambda x: x[0], reverse=True) 

3456 

3457 matching_knowledge_bases = [] 

3458 for distance, knowledge_base_id in sorted_results: 

3459 knowledge_base = await Knowledges.get_knowledge_by_id(knowledge_base_id) 

3460 if knowledge_base: 

3461 matching_knowledge_bases.append( 

3462 { 

3463 'id': knowledge_base.id, 

3464 'name': knowledge_base.name, 

3465 'description': knowledge_base.description or '', 

3466 'similarity': round(distance, 4), 

3467 } 

3468 ) 

3469 

3470 return JSONCodec.dumps(matching_knowledge_bases, ensure_ascii=False) 

3471 

3472 except Exception as e: 

3473 log.exception(f'query_knowledge_bases error: {e}') 

3474 return JSONCodec.dumps({'error': str(e)}) 

3475 

3476 

3477# ============================================================================= 

3478# SKILLS TOOLS 

3479# ============================================================================= 

3480 

3481 

3482async def view_skill( 

3483 id: str, 

3484 __request__: Request = None, 

3485 __user__: dict = None, 

3486 __metadata__: dict = None, 

3487) -> str: 

3488 """ 

3489 Load the full instructions of a skill by its id from the available skills manifest. 

3490 Use this when you need detailed instructions for a skill listed in <available_skills>. 

3491 

3492 :param id: The id of the skill to load (as shown in the manifest) 

3493 :return: The full skill instructions as markdown content 

3494 """ 

3495 if __request__ is None: 

3496 return JSONCodec.dumps({'error': 'Request context not available'}) 

3497 

3498 if not __user__: 

3499 return JSONCodec.dumps({'error': 'User context not available'}) 

3500 

3501 try: 

3502 terminal_skill_prefix = 'terminal:' 

3503 if isinstance(id, str) and id.startswith(terminal_skill_prefix): 

3504 from open_webui.utils.terminals import get_terminal_skill 

3505 

3506 skill_name = unquote(id.removeprefix(terminal_skill_prefix)) 

3507 skill = await get_terminal_skill(__request__, __user__, __metadata__ or {}, skill_name) 

3508 if not skill: 

3509 return JSONCodec.dumps({'error': f"Skill '{id}' not found"}) 

3510 return JSONCodec.dumps(skill, ensure_ascii=False) 

3511 

3512 from open_webui.models.access_grants import AccessGrants 

3513 from open_webui.models.skills import Skills 

3514 

3515 user_id = __user__.get('id') 

3516 

3517 # Direct DB lookup by id (case-insensitive since IDs are stored lowercase) 

3518 skill = await Skills.get_skill_by_id(id.lower()) 

3519 

3520 if not skill or not skill.is_active: 

3521 return JSONCodec.dumps({'error': f"Skill '{id}' not found"}) 

3522 

3523 # Check user access 

3524 user_role = __user__.get('role', 'user') 

3525 if user_role != 'admin' and skill.user_id != user_id: 

3526 user_group_ids = [group.id for group in await Groups.get_groups_by_member_id(user_id)] 

3527 if not await AccessGrants.has_access( 

3528 user_id=user_id, 

3529 resource_type='skill', 

3530 resource_id=skill.id, 

3531 permission='read', 

3532 user_group_ids=set(user_group_ids), 

3533 ): 

3534 return JSONCodec.dumps({'error': 'Access denied'}) 

3535 

3536 return JSONCodec.dumps( 

3537 { 

3538 'name': skill.name, 

3539 'content': skill.content, 

3540 }, 

3541 ensure_ascii=False, 

3542 ) 

3543 except Exception as e: 

3544 log.exception(f'view_skill error: {e}') 

3545 return JSONCodec.dumps({'error': str(e)}) 

3546 

3547 

3548# ============================================================================= 

3549# TASK MANAGEMENT TOOLS 

3550# ============================================================================= 

3551 

3552from typing import Literal 

3553 

3554from pydantic import BaseModel, Field 

3555 

3556VALID_TASK_STATUSES = {'pending', 'in_progress', 'completed', 'cancelled'} 

3557 

3558 

3559class TaskItem(BaseModel): 

3560 id: Optional[str] = Field(None, description='Unique identifier for the task. Auto-generated if omitted.') 

3561 content: str = Field(..., description='Task description.') 

3562 status: Literal['pending', 'in_progress', 'completed', 'cancelled'] = Field('pending', description='Task status.') 

3563 

3564 

3565def _task_summary(all_tasks: list[dict]) -> dict: 

3566 """Build summary counts for a task list.""" 

3567 pending = sum(1 for t in all_tasks if t['status'] == 'pending') 

3568 in_progress = sum(1 for t in all_tasks if t['status'] == 'in_progress') 

3569 completed = sum(1 for t in all_tasks if t['status'] == 'completed') 

3570 cancelled = sum(1 for t in all_tasks if t['status'] == 'cancelled') 

3571 return { 

3572 'total': len(all_tasks), 

3573 'pending': pending, 

3574 'in_progress': in_progress, 

3575 'completed': completed, 

3576 'cancelled': cancelled, 

3577 } 

3578 

3579 

3580async def _emit_tasks(event_emitter, all_tasks: list[dict]): 

3581 """Persist task state to the UI.""" 

3582 if event_emitter: 

3583 await event_emitter( 

3584 { 

3585 'type': 'chat:message:tasks', 

3586 'data': { 

3587 'tasks': all_tasks, 

3588 }, 

3589 } 

3590 ) 

3591 

3592 

3593async def create_tasks( 

3594 tasks: list[TaskItem], 

3595 __chat_id__: str = None, 

3596 __message_id__: str = None, 

3597 __event_emitter__: callable = None, 

3598 __request__: Request = None, 

3599 __user__: dict = None, 

3600) -> str: 

3601 """ 

3602 Create a visible task checklist for multi-step work so progress can be shown in chat. 

3603 

3604 :param tasks: List of task items. Each item: content (string, required), status (pending|in_progress|completed|cancelled, default pending), id (optional, auto-generated). 

3605 :return: JSON with the full task list and summary counts 

3606 """ 

3607 if not is_saved_chat_id(__chat_id__): 

3608 return JSONCodec.dumps({'error': 'Saved chat context not available'}) 

3609 

3610 try: 

3611 all_tasks = [] 

3612 for idx, task in enumerate(tasks): 

3613 if hasattr(task, 'model_dump'): 

3614 d = task.model_dump(exclude_none=True) 

3615 elif isinstance(task, dict): 

3616 d = task 

3617 else: 

3618 d = dict(task) 

3619 

3620 content = str(d.get('content', '')).strip() 

3621 if not content: 

3622 continue 

3623 

3624 item_id = str(d.get('id', '') or '').strip() or str(idx + 1) 

3625 status = str(d.get('status', 'pending')).strip().lower() 

3626 if status not in VALID_TASK_STATUSES: 

3627 status = 'pending' 

3628 

3629 all_tasks.append({'id': item_id, 'content': content, 'status': status}) 

3630 

3631 await Chats.update_chat_tasks_by_id(__chat_id__, all_tasks) 

3632 await _emit_tasks(__event_emitter__, all_tasks) 

3633 

3634 return JSONCodec.dumps( 

3635 {'tasks': all_tasks, 'summary': _task_summary(all_tasks)}, 

3636 ensure_ascii=False, 

3637 ) 

3638 except Exception as e: 

3639 log.exception(f'tasks error: {e}') 

3640 return JSONCodec.dumps({'error': str(e)}) 

3641 

3642 

3643async def update_task( 

3644 id: str, 

3645 status: str = 'completed', 

3646 __chat_id__: str = None, 

3647 __message_id__: str = None, 

3648 __event_emitter__: callable = None, 

3649 __request__: Request = None, 

3650 __user__: dict = None, 

3651) -> str: 

3652 """ 

3653 Mark a single visible task item as completed, in_progress, pending, or cancelled. 

3654 

3655 :param id: The task ID to update 

3656 :param status: New status: completed, in_progress, pending, or cancelled (default: completed) 

3657 :return: JSON with the updated task list and summary counts 

3658 """ 

3659 if not is_saved_chat_id(__chat_id__): 

3660 return JSONCodec.dumps({'error': 'Saved chat context not available'}) 

3661 

3662 try: 

3663 status = status.strip().lower() 

3664 if status not in VALID_TASK_STATUSES: 

3665 return JSONCodec.dumps( 

3666 {'error': f'Invalid status: {status}. Must be one of: {", ".join(sorted(VALID_TASK_STATUSES))}'} 

3667 ) 

3668 

3669 all_tasks = await Chats.get_chat_tasks_by_id(__chat_id__) 

3670 

3671 found = False 

3672 for task in all_tasks: 

3673 if task['id'] == id: 

3674 task['status'] = status 

3675 found = True 

3676 break 

3677 

3678 if not found: 

3679 return JSONCodec.dumps({'error': f'Task with id "{id}" not found'}) 

3680 

3681 await Chats.update_chat_tasks_by_id(__chat_id__, all_tasks) 

3682 await _emit_tasks(__event_emitter__, all_tasks) 

3683 

3684 return JSONCodec.dumps( 

3685 {'tasks': all_tasks, 'summary': _task_summary(all_tasks)}, 

3686 ensure_ascii=False, 

3687 ) 

3688 except Exception as e: 

3689 log.exception(f'update_task_status error: {e}') 

3690 return JSONCodec.dumps({'error': str(e)}) 

3691 

3692 

3693# ============================================================================= 

3694# AUTOMATION TOOLS 

3695# ============================================================================= 

3696 

3697 

3698async def _validate_owned_automation_folder(user_id: str, folder_id: Optional[str]) -> Optional[str]: 

3699 if not folder_id: 

3700 return None 

3701 from open_webui.models.folders import Folders 

3702 

3703 folder = await Folders.get_folder_by_id_and_user_id(folder_id, user_id) 

3704 if not folder: 

3705 raise ValueError('Folder not found') 

3706 return folder.id 

3707 

3708 

3709async def create_automation( 

3710 name: str, 

3711 prompt: str, 

3712 rrule: str, 

3713 folder_id: Optional[str] = None, 

3714 __request__: Request = None, 

3715 __user__: dict = None, 

3716 __metadata__: dict = None, 

3717) -> str: 

3718 """ 

3719 Create a scheduled automation that runs a prompt on a recurring or one-time schedule. 

3720 Use this when the user wants to schedule a task to run automatically. 

3721 The automation will use the current chat model. 

3722 

3723 The rrule parameter must be a valid iCalendar RRULE string. Common examples: 

3724 - Every day at 9am: "DTSTART:20250101T090000\\nRRULE:FREQ=DAILY" 

3725 - Every Monday at 8am: "DTSTART:20250106T080000\\nRRULE:FREQ=WEEKLY;BYDAY=MO" 

3726 - Every hour: "RRULE:FREQ=HOURLY;INTERVAL=1" 

3727 - Every 30 minutes: "RRULE:FREQ=MINUTELY;INTERVAL=30" 

3728 - Once at a specific time: "DTSTART:20250415T140000\\nRRULE:FREQ=DAILY;COUNT=1" 

3729 - First day of every month: "DTSTART:20250101T090000\\nRRULE:FREQ=MONTHLY;BYMONTHDAY=1" 

3730 

3731 The DTSTART time should reflect the desired execution time. Use COUNT=1 for one-time automations. 

3732 

3733 :param name: A short descriptive name for the automation 

3734 :param prompt: The prompt/instructions to execute on each run 

3735 :param rrule: An iCalendar RRULE string defining the schedule 

3736 :param folder_id: Optional owner-owned folder ID for generated chats 

3737 :return: JSON with the created automation details including id, next scheduled runs 

3738 """ 

3739 if __request__ is None: 

3740 return JSONCodec.dumps({'error': 'Request context not available'}) 

3741 

3742 if not __user__: 

3743 return JSONCodec.dumps({'error': 'User context not available'}) 

3744 

3745 try: 

3746 from open_webui.models.automations import AutomationData, AutomationForm, AutomationTarget, Automations 

3747 from open_webui.models.users import Users 

3748 from open_webui.routers.automations import check_automation_limits 

3749 from open_webui.utils.automations import next_n_runs_ns, next_run_ns, validate_rrule 

3750 

3751 user_id = __user__.get('id') 

3752 user = await Users.get_user_by_id(user_id) 

3753 if not user: 

3754 return JSONCodec.dumps({'error': 'User not found'}) 

3755 

3756 # Fall back to model dict ID since __metadata__ may predate model_id assignment 

3757 metadata = __metadata__ or {} 

3758 model_id = metadata.get('model_id') or ( 

3759 metadata.get('model', {}).get('id') if isinstance(metadata.get('model'), dict) else None 

3760 ) 

3761 if not model_id: 

3762 return JSONCodec.dumps({'error': 'Could not detect current model'}) 

3763 

3764 try: 

3765 folder_id = await _validate_owned_automation_folder(user_id, folder_id) 

3766 except ValueError as e: 

3767 return JSONCodec.dumps({'error': str(e)}) 

3768 

3769 # Validate the RRULE 

3770 try: 

3771 await validate_rrule(rrule, tz=user.timezone) 

3772 except ValueError as e: 

3773 return JSONCodec.dumps({'error': f'Invalid schedule: {e}'}) 

3774 

3775 try: 

3776 await check_automation_limits(__request__, user, rrule, None, is_create=True) 

3777 except HTTPException as e: 

3778 return JSONCodec.dumps({'error': e.detail}) 

3779 

3780 tz = user.timezone 

3781 form = AutomationForm( 

3782 name=name, 

3783 folder_id=folder_id, 

3784 data=AutomationData( 

3785 prompt=prompt, 

3786 model_id=model_id, 

3787 rrule=rrule, 

3788 target=( 

3789 AutomationTarget(type='channel', channel_id=metadata.get('chat_id', '').removeprefix('channel:')) 

3790 if metadata.get('chat_id', '').startswith('channel:') 

3791 else None 

3792 ), 

3793 ), 

3794 is_active=True, 

3795 ) 

3796 

3797 automation = await Automations.insert(user_id, form, await next_run_ns(rrule, tz=tz)) 

3798 

3799 return JSONCodec.dumps( 

3800 { 

3801 'status': 'success', 

3802 'id': automation.id, 

3803 'name': automation.name, 

3804 'folder_id': automation.folder_id, 

3805 'model_id': model_id, 

3806 'target': automation.data.get('target'), 

3807 'is_active': automation.is_active, 

3808 'next_runs': await next_n_runs_ns(rrule, tz=tz), 

3809 }, 

3810 ensure_ascii=False, 

3811 ) 

3812 except Exception as e: 

3813 log.exception(f'create_automation error: {e}') 

3814 return JSONCodec.dumps({'error': str(e)}) 

3815 

3816 

3817async def update_automation( 

3818 automation_id: str, 

3819 name: Optional[str] = None, 

3820 prompt: Optional[str] = None, 

3821 rrule: Optional[str] = None, 

3822 model_id: Optional[str] = None, 

3823 folder_id: Optional[str] = '', 

3824 __request__: Request = None, 

3825 __user__: dict = None, 

3826) -> str: 

3827 """ 

3828 Update an existing automation. Only the provided fields are changed; omitted fields stay the same. 

3829 

3830 :param automation_id: The ID of the automation to update 

3831 :param name: New name for the automation (optional) 

3832 :param prompt: New prompt/instructions (optional) 

3833 :param rrule: New iCalendar RRULE schedule string (optional). See create_automation for format examples. 

3834 :param model_id: New model ID to use (optional); blank values are ignored 

3835 :param folder_id: New owner-owned folder ID (optional); omit or pass blank to keep unchanged, pass null to clear 

3836 :return: JSON with the updated automation details 

3837 """ 

3838 if __request__ is None: 

3839 return JSONCodec.dumps({'error': 'Request context not available'}) 

3840 

3841 if not __user__: 

3842 return JSONCodec.dumps({'error': 'User context not available'}) 

3843 

3844 try: 

3845 from open_webui.models.automations import AutomationData, AutomationForm, AutomationTarget, Automations 

3846 from open_webui.models.users import Users 

3847 from open_webui.routers.automations import check_automation_limits 

3848 from open_webui.utils.automations import next_n_runs_ns, next_run_ns, validate_rrule 

3849 

3850 user_id = __user__.get('id') 

3851 user = await Users.get_user_by_id(user_id) 

3852 if not user: 

3853 return JSONCodec.dumps({'error': 'User not found'}) 

3854 

3855 automation = await Automations.get_by_id(automation_id) 

3856 if not automation: 

3857 return JSONCodec.dumps({'error': 'Automation not found'}) 

3858 if automation.user_id != user_id: 

3859 return JSONCodec.dumps({'error': 'Access denied'}) 

3860 

3861 # Merge provided fields with existing values 

3862 new_name = name if name is not None else automation.name 

3863 new_prompt = prompt if prompt is not None else automation.data.get('prompt', '') 

3864 new_model_id = model_id.strip() if model_id and model_id.strip() else automation.data.get('model_id', '') 

3865 new_rrule = rrule if rrule is not None else automation.data.get('rrule', '') 

3866 if folder_id is None: 

3867 new_folder_id = None 

3868 elif not folder_id.strip(): 

3869 new_folder_id = automation.folder_id 

3870 else: 

3871 try: 

3872 new_folder_id = await _validate_owned_automation_folder(user_id, folder_id.strip()) 

3873 except ValueError as e: 

3874 return JSONCodec.dumps({'error': str(e)}) 

3875 

3876 # Validate RRULE if changed 

3877 if rrule is not None: 

3878 try: 

3879 await validate_rrule(new_rrule, tz=user.timezone) 

3880 except ValueError as e: 

3881 return JSONCodec.dumps({'error': f'Invalid schedule: {e}'}) 

3882 

3883 try: 

3884 await check_automation_limits(__request__, user, new_rrule, None) 

3885 except HTTPException as e: 

3886 return JSONCodec.dumps({'error': e.detail}) 

3887 

3888 tz = user.timezone 

3889 form = AutomationForm( 

3890 name=new_name, 

3891 folder_id=new_folder_id, 

3892 data=AutomationData( 

3893 prompt=new_prompt, 

3894 model_id=new_model_id, 

3895 rrule=new_rrule, 

3896 target=AutomationTarget(**automation.data['target']) if automation.data.get('target') else None, 

3897 ), 

3898 is_active=automation.is_active, 

3899 ) 

3900 

3901 updated = await Automations.update_by_id(automation_id, form, await next_run_ns(new_rrule, tz=tz)) 

3902 

3903 return JSONCodec.dumps( 

3904 { 

3905 'status': 'success', 

3906 'id': updated.id, 

3907 'name': updated.name, 

3908 'folder_id': updated.folder_id, 

3909 'model_id': new_model_id, 

3910 'target': updated.data.get('target'), 

3911 'is_active': updated.is_active, 

3912 'next_runs': await next_n_runs_ns(new_rrule, tz=tz), 

3913 }, 

3914 ensure_ascii=False, 

3915 ) 

3916 except Exception as e: 

3917 log.exception(f'update_automation error: {e}') 

3918 return JSONCodec.dumps({'error': str(e)}) 

3919 

3920 

3921async def list_automations( 

3922 status: Optional[str] = None, 

3923 folder_id: Optional[str] = None, 

3924 count: int = 10, 

3925 __request__: Request = None, 

3926 __user__: dict = None, 

3927) -> str: 

3928 """ 

3929 List the user's scheduled automations. 

3930 

3931 :param status: Filter by status: "active", "paused", or omit for all 

3932 :param folder_id: Optional owner-owned folder ID filter; pass an empty string to clear the folder filter 

3933 :param count: Maximum number of automations to return (default: 10) 

3934 :return: JSON list of automations with id, name, prompt snippet, schedule, status, and next runs 

3935 """ 

3936 if __request__ is None: 

3937 return JSONCodec.dumps({'error': 'Request context not available'}) 

3938 

3939 if not __user__: 

3940 return JSONCodec.dumps({'error': 'User context not available'}) 

3941 

3942 try: 

3943 from open_webui.models.automations import Automations 

3944 from open_webui.models.users import Users 

3945 from open_webui.utils.automations import next_n_runs_ns 

3946 

3947 user_id = __user__.get('id') 

3948 user = await Users.get_user_by_id(user_id) 

3949 if folder_id: 

3950 try: 

3951 folder_id = await _validate_owned_automation_folder(user_id, folder_id) 

3952 except ValueError as e: 

3953 return JSONCodec.dumps({'error': str(e)}) 

3954 

3955 result = await Automations.search_automations( 

3956 user_id=user_id, 

3957 status=status, 

3958 folder_id=folder_id, 

3959 skip=0, 

3960 limit=count, 

3961 ) 

3962 

3963 automations = [] 

3964 for item in result.items: 

3965 rrule = item.data.get('rrule', '') 

3966 prompt_text = item.data.get('prompt', '') 

3967 snippet = prompt_text[:100] + ('...' if len(prompt_text) > 100 else '') 

3968 

3969 automations.append( 

3970 { 

3971 'id': item.id, 

3972 'name': item.name, 

3973 'folder_id': item.folder_id, 

3974 'prompt_snippet': snippet, 

3975 'model_id': item.data.get('model_id', ''), 

3976 'target': item.data.get('target'), 

3977 'rrule': rrule, 

3978 'is_active': item.is_active, 

3979 'last_run_at': item.last_run_at, 

3980 'next_runs': await next_n_runs_ns(rrule, tz=user.timezone if user else None), 

3981 } 

3982 ) 

3983 

3984 return JSONCodec.dumps( 

3985 {'automations': automations, 'total': result.total}, 

3986 ensure_ascii=False, 

3987 ) 

3988 except Exception as e: 

3989 log.exception(f'list_automations error: {e}') 

3990 return JSONCodec.dumps({'error': str(e)}) 

3991 

3992 

3993async def toggle_automation( 

3994 automation_id: str, 

3995 __request__: Request = None, 

3996 __user__: dict = None, 

3997) -> str: 

3998 """ 

3999 Pause or resume a scheduled automation. If active, it will be paused. If paused, it will be resumed. 

4000 

4001 :param automation_id: The ID of the automation to toggle 

4002 :return: JSON with the updated automation status 

4003 """ 

4004 if __request__ is None: 

4005 return JSONCodec.dumps({'error': 'Request context not available'}) 

4006 

4007 if not __user__: 

4008 return JSONCodec.dumps({'error': 'User context not available'}) 

4009 

4010 try: 

4011 from open_webui.models.automations import Automations 

4012 from open_webui.models.users import Users 

4013 from open_webui.utils.automations import next_run_ns 

4014 

4015 user_id = __user__.get('id') 

4016 user = await Users.get_user_by_id(user_id) 

4017 

4018 automation = await Automations.get_by_id(automation_id) 

4019 if not automation: 

4020 return JSONCodec.dumps({'error': 'Automation not found'}) 

4021 if automation.user_id != user_id: 

4022 return JSONCodec.dumps({'error': 'Access denied'}) 

4023 

4024 rrule = automation.data.get('rrule', '') 

4025 toggled = await Automations.toggle( 

4026 automation_id, 

4027 await next_run_ns(rrule, tz=user.timezone if user else None), 

4028 ) 

4029 

4030 return JSONCodec.dumps( 

4031 { 

4032 'status': 'success', 

4033 'id': toggled.id, 

4034 'name': toggled.name, 

4035 'is_active': toggled.is_active, 

4036 }, 

4037 ensure_ascii=False, 

4038 ) 

4039 except Exception as e: 

4040 log.exception(f'toggle_automation error: {e}') 

4041 return JSONCodec.dumps({'error': str(e)}) 

4042 

4043 

4044async def delete_automation( 

4045 automation_id: str, 

4046 __request__: Request = None, 

4047 __user__: dict = None, 

4048) -> str: 

4049 """ 

4050 Delete a scheduled automation and all its run history. 

4051 

4052 :param automation_id: The ID of the automation to delete 

4053 :return: JSON confirming the automation was deleted 

4054 """ 

4055 if __request__ is None: 

4056 return JSONCodec.dumps({'error': 'Request context not available'}) 

4057 

4058 if not __user__: 

4059 return JSONCodec.dumps({'error': 'User context not available'}) 

4060 

4061 try: 

4062 from open_webui.models.automations import AutomationRuns, Automations 

4063 

4064 user_id = __user__.get('id') 

4065 

4066 automation = await Automations.get_by_id(automation_id) 

4067 if not automation: 

4068 return JSONCodec.dumps({'error': 'Automation not found'}) 

4069 if automation.user_id != user_id: 

4070 return JSONCodec.dumps({'error': 'Access denied'}) 

4071 

4072 name = automation.name 

4073 await AutomationRuns.delete_by_automation(automation_id) 

4074 await Automations.delete(automation_id) 

4075 

4076 return JSONCodec.dumps( 

4077 { 

4078 'status': 'success', 

4079 'message': f'Automation "{name}" deleted', 

4080 }, 

4081 ensure_ascii=False, 

4082 ) 

4083 except Exception as e: 

4084 log.exception(f'delete_automation error: {e}') 

4085 return JSONCodec.dumps({'error': str(e)}) 

4086 

4087 

4088# ============================================================================= 

4089# CALENDAR TOOLS 

4090# ============================================================================= 

4091 

4092 

4093MAX_CALENDAR_RANGE_END_NS = 2**63 - 1 

4094 

4095 

4096def _get_user_tz(user_dict: dict): 

4097 """Get the user's timezone as a ZoneInfo, falling back to UTC.""" 

4098 from zoneinfo import ZoneInfo 

4099 

4100 tz_name = None 

4101 if user_dict: 

4102 tz_name = user_dict.get('timezone') 

4103 if tz_name: 

4104 try: 

4105 return ZoneInfo(tz_name) 

4106 except Exception: 

4107 pass 

4108 return ZoneInfo('UTC') 

4109 

4110 

4111def _dt_to_ns(dt_str: str, tz) -> int: 

4112 """Convert a datetime string to nanoseconds since epoch, interpreting in the given timezone.""" 

4113 from datetime import datetime 

4114 

4115 dt = datetime.fromisoformat(dt_str) 

4116 # If naive (no timezone info), localize to user's timezone 

4117 if dt.tzinfo is None: 

4118 dt = dt.replace(tzinfo=tz) 

4119 return int(dt.timestamp() * 1_000) * 1_000_000 

4120 

4121 

4122def _ns_to_dt(ns: int, tz) -> str: 

4123 """Convert nanoseconds since epoch to a datetime string in the given timezone.""" 

4124 from datetime import datetime 

4125 

4126 seconds = ns / 1_000_000_000 

4127 dt = datetime.fromtimestamp(seconds, tz=tz) 

4128 return dt.strftime('%Y-%m-%d %H:%M') 

4129 

4130 

4131def _event_to_dict(event, tz) -> dict: 

4132 """Convert a calendar event model to a human-friendly dict with local timestamps.""" 

4133 alert_minutes = None 

4134 if event.meta and 'alert_minutes' in event.meta: 

4135 alert_minutes = event.meta['alert_minutes'] 

4136 return { 

4137 'id': event.id, 

4138 'calendar_id': event.calendar_id, 

4139 'title': event.title, 

4140 'description': event.description or '', 

4141 'start': _ns_to_dt(event.start_at, tz), 

4142 'end': _ns_to_dt(event.end_at, tz) if event.end_at else None, 

4143 'all_day': event.all_day, 

4144 'location': event.location or '', 

4145 'reminder_minutes': alert_minutes if alert_minutes is not None else 10, 

4146 'color': event.color, 

4147 'is_cancelled': event.is_cancelled, 

4148 } 

4149 

4150 

4151async def search_calendar_events( 

4152 query: Optional[str] = None, 

4153 start: Optional[str] = None, 

4154 end: Optional[str] = None, 

4155 count: int = 10, 

4156 __request__: Request = None, 

4157 __user__: dict = None, 

4158) -> str: 

4159 """ 

4160 Search calendar events, reminders, and scheduled items by text and/or date range. 

4161 Helpful for finding upcoming events, reminders, or schedule items. 

4162 

4163 :param query: Search text to match against event title, description, or location (optional) 

4164 :param start: Only return events starting at or after this datetime, e.g. "2026-04-20 00:00" (optional) 

4165 :param end: Only return events starting before this datetime, e.g. "2026-04-27 00:00" (optional) 

4166 :param count: Maximum number of events to return (default: 10) 

4167 :return: JSON list of matching events with id, title, description, start, end, calendar_id, location 

4168 """ 

4169 if __request__ is None: 

4170 return JSONCodec.dumps({'error': 'Request context not available'}) 

4171 

4172 if not __user__: 

4173 return JSONCodec.dumps({'error': 'User context not available'}) 

4174 

4175 try: 

4176 from open_webui.models.calendar import CalendarEvents 

4177 

4178 user_id = __user__.get('id') 

4179 tz = _get_user_tz(__user__) 

4180 

4181 if isinstance(count, str): 

4182 try: 

4183 count = int(count) 

4184 except ValueError: 

4185 count = 10 

4186 

4187 if start or end: 

4188 # Date range query — use get_events_by_range 

4189 try: 

4190 start_ns = _dt_to_ns(start, tz) if start else 0 

4191 except (ValueError, TypeError) as e: 

4192 return JSONCodec.dumps({'error': f'Invalid start datetime: {e}'}) 

4193 

4194 try: 

4195 end_ns = _dt_to_ns(end, tz) if end else MAX_CALENDAR_RANGE_END_NS 

4196 except (ValueError, TypeError) as e: 

4197 return JSONCodec.dumps({'error': f'Invalid end datetime: {e}'}) 

4198 

4199 items = await CalendarEvents.get_events_by_range( 

4200 user_id=user_id, 

4201 start=start_ns, 

4202 end=end_ns, 

4203 ) 

4204 

4205 # Apply text filter if query is also provided 

4206 if query: 

4207 q = query.lower() 

4208 items = [ 

4209 e 

4210 for e in items 

4211 if q in (e.title or '').lower() 

4212 or q in (e.description or '').lower() 

4213 or q in (e.location or '').lower() 

4214 ] 

4215 

4216 events = [_event_to_dict(item, tz) for item in items[:count]] 

4217 return JSONCodec.dumps( 

4218 {'events': events, 'total': len(items)}, 

4219 ensure_ascii=False, 

4220 ) 

4221 else: 

4222 # Text-only search 

4223 result = await CalendarEvents.search_events( 

4224 user_id=user_id, 

4225 query=query, 

4226 skip=0, 

4227 limit=count, 

4228 ) 

4229 

4230 events = [_event_to_dict(item, tz) for item in result.items] 

4231 return JSONCodec.dumps( 

4232 {'events': events, 'total': result.total}, 

4233 ensure_ascii=False, 

4234 ) 

4235 except Exception as e: 

4236 log.exception(f'search_calendar_events error: {e}') 

4237 return JSONCodec.dumps({'error': str(e)}) 

4238 

4239 

4240async def create_calendar_event( 

4241 title: str, 

4242 start: str, 

4243 end: Optional[str] = None, 

4244 description: Optional[str] = None, 

4245 calendar_id: Optional[str] = None, 

4246 all_day: bool = False, 

4247 location: Optional[str] = None, 

4248 reminder_minutes: Optional[int] = None, 

4249 __request__: Request = None, 

4250 __user__: dict = None, 

4251) -> str: 

4252 """ 

4253 Create a calendar event, reminder, or alarm. Use this when the user wants to 

4254 schedule an event, set a reminder, create an alarm, or says things like 

4255 "remind me", "don't let me forget", "notify me at", or "add to my calendar". 

4256 For simple reminders, omit end/location/all_day and set reminder_minutes to 0. 

4257 

4258 :param title: Event or reminder title (e.g. "Team standup", "Take medicine", "Call mom") 

4259 :param start: Start datetime in the user's local time (e.g. "2026-04-20 09:00") 

4260 :param end: End datetime in the user's local time (optional — omit for reminders or point-in-time events) 

4261 :param description: Event description or notes (optional) 

4262 :param calendar_id: Target calendar ID (optional, uses default calendar if omitted) 

4263 :param all_day: Whether this is an all-day event (default: false) 

4264 :param location: Event location (optional) 

4265 :param reminder_minutes: Minutes before the event to send a notification (optional, default: 10). Use 0 for "at time of event", -1 for no notification. 

4266 :return: JSON with the created event details including id 

4267 """ 

4268 if __request__ is None: 

4269 return JSONCodec.dumps({'error': 'Request context not available'}) 

4270 

4271 if not __user__: 

4272 return JSONCodec.dumps({'error': 'User context not available'}) 

4273 

4274 try: 

4275 from open_webui.models.calendar import CalendarEventForm, CalendarEvents, Calendars 

4276 

4277 user_id = __user__.get('id') 

4278 

4279 # Resolve calendar_id: use provided, or fall back to default 

4280 if not calendar_id: 

4281 calendars = await Calendars.get_calendars_by_user(user_id) 

4282 default_cal = next((c for c in calendars if c.is_default), None) 

4283 if not default_cal and calendars: 

4284 default_cal = calendars[0] 

4285 if not default_cal: 

4286 return JSONCodec.dumps({'error': 'No calendars found. Cannot create event.'}) 

4287 calendar_id = default_cal.id 

4288 

4289 # Verify access 

4290 cal = await Calendars.get_calendar_by_id(calendar_id) 

4291 if not cal: 

4292 return JSONCodec.dumps({'error': 'Calendar not found'}) 

4293 if cal.user_id != user_id and __user__.get('role') != 'admin': 

4294 from open_webui.models.access_grants import AccessGrants 

4295 from open_webui.models.groups import Groups 

4296 

4297 user_group_ids = [g.id for g in await Groups.get_groups_by_member_id(user_id)] 

4298 if not await AccessGrants.has_access( 

4299 user_id=user_id, 

4300 resource_type='calendar', 

4301 resource_id=cal.id, 

4302 permission='write', 

4303 user_group_ids=set(user_group_ids), 

4304 ): 

4305 return JSONCodec.dumps({'error': 'Access denied to this calendar'}) 

4306 

4307 # Coerce boolean from LLM 

4308 if isinstance(all_day, str): 

4309 all_day = all_day.lower() in ('true', '1', 'yes') 

4310 

4311 # Convert datetime strings to nanoseconds using user's timezone 

4312 tz = _get_user_tz(__user__) 

4313 try: 

4314 start_ns = _dt_to_ns(start, tz) 

4315 except (ValueError, TypeError) as e: 

4316 return JSONCodec.dumps({'error': f'Invalid start datetime: {e}. Use format like "2026-04-20 09:00"'}) 

4317 

4318 end_ns = None 

4319 if end: 

4320 try: 

4321 end_ns = _dt_to_ns(end, tz) 

4322 except (ValueError, TypeError) as e: 

4323 return JSONCodec.dumps({'error': f'Invalid end datetime: {e}. Use format like "2026-04-20 10:00"'}) 

4324 elif not all_day: 

4325 # Default to 1 hour duration 

4326 end_ns = start_ns + 3_600_000_000_000 

4327 

4328 # Build meta with reminder setting 

4329 meta = {} 

4330 if reminder_minutes is not None: 

4331 if isinstance(reminder_minutes, str): 

4332 try: 

4333 reminder_minutes = int(reminder_minutes) 

4334 except ValueError: 

4335 reminder_minutes = 10 

4336 meta['alert_minutes'] = reminder_minutes 

4337 else: 

4338 meta['alert_minutes'] = 10 

4339 

4340 form = CalendarEventForm( 

4341 calendar_id=calendar_id, 

4342 title=title, 

4343 description=description, 

4344 start_at=start_ns, 

4345 end_at=end_ns, 

4346 all_day=all_day, 

4347 location=location, 

4348 meta=meta, 

4349 ) 

4350 

4351 event = await CalendarEvents.insert_new_event(user_id, form) 

4352 if not event: 

4353 return JSONCodec.dumps({'error': 'Failed to create event'}) 

4354 

4355 return JSONCodec.dumps( 

4356 { 

4357 'status': 'success', 

4358 **_event_to_dict(event, tz), 

4359 }, 

4360 ensure_ascii=False, 

4361 ) 

4362 except Exception as e: 

4363 log.exception(f'create_calendar_event error: {e}') 

4364 return JSONCodec.dumps({'error': str(e)}) 

4365 

4366 

4367async def update_calendar_event( 

4368 event_id: str, 

4369 title: Optional[str] = None, 

4370 description: Optional[str] = None, 

4371 start: Optional[str] = None, 

4372 end: Optional[str] = None, 

4373 all_day: Optional[bool] = None, 

4374 location: Optional[str] = None, 

4375 is_cancelled: Optional[bool] = None, 

4376 reminder_minutes: Optional[int] = None, 

4377 __request__: Request = None, 

4378 __user__: dict = None, 

4379) -> str: 

4380 """ 

4381 Update an existing calendar event. Only provided fields are changed; 

4382 omitted fields stay the same. 

4383 

4384 :param event_id: The ID of the event to update 

4385 :param title: New event title (optional) 

4386 :param description: New event description (optional) 

4387 :param start: New start datetime string in your local time, e.g. "2026-04-20 09:00" (optional) 

4388 :param end: New end datetime string in your local time (optional) 

4389 :param all_day: Whether this is an all-day event (optional) 

4390 :param location: New event location (optional) 

4391 :param is_cancelled: Set to true to cancel the event (optional) 

4392 :param reminder_minutes: Minutes before the event to send a reminder notification (optional). Use 0 for "at time of event", -1 for no reminder. Accepts any positive integer for custom timing (e.g. 120 for 2 hours before). 

4393 :return: JSON with the updated event details 

4394 """ 

4395 if __request__ is None: 

4396 return JSONCodec.dumps({'error': 'Request context not available'}) 

4397 

4398 if not __user__: 

4399 return JSONCodec.dumps({'error': 'User context not available'}) 

4400 

4401 try: 

4402 from open_webui.models.access_grants import AccessGrants 

4403 from open_webui.models.calendar import CalendarEvents, CalendarEventUpdateForm, Calendars 

4404 from open_webui.models.groups import Groups 

4405 

4406 user_id = __user__.get('id') 

4407 

4408 event = await CalendarEvents.get_event_by_id(event_id) 

4409 if not event: 

4410 return JSONCodec.dumps({'error': 'Event not found'}) 

4411 

4412 # Check write access to the event's calendar 

4413 if event.user_id != user_id and __user__.get('role') != 'admin': 

4414 cal = await Calendars.get_calendar_by_id(event.calendar_id) 

4415 if not cal: 

4416 return JSONCodec.dumps({'error': 'Access denied'}) 

4417 user_group_ids = [g.id for g in await Groups.get_groups_by_member_id(user_id)] 

4418 if not await AccessGrants.has_access( 

4419 user_id=user_id, 

4420 resource_type='calendar', 

4421 resource_id=cal.id, 

4422 permission='write', 

4423 user_group_ids=set(user_group_ids), 

4424 ): 

4425 return JSONCodec.dumps({'error': 'Access denied'}) 

4426 

4427 # Coerce boolean strings from LLM 

4428 if isinstance(all_day, str): 

4429 all_day = all_day.lower() in ('true', '1', 'yes') 

4430 if isinstance(is_cancelled, str): 

4431 is_cancelled = is_cancelled.lower() in ('true', '1', 'yes') 

4432 

4433 # Convert datetime strings to nanoseconds using user's timezone 

4434 tz = _get_user_tz(__user__) 

4435 start_ns = None 

4436 if start is not None: 

4437 try: 

4438 start_ns = _dt_to_ns(start, tz) 

4439 except (ValueError, TypeError) as e: 

4440 return JSONCodec.dumps({'error': f'Invalid start datetime: {e}'}) 

4441 

4442 end_ns = None 

4443 if end is not None: 

4444 try: 

4445 end_ns = _dt_to_ns(end, tz) 

4446 except (ValueError, TypeError) as e: 

4447 return JSONCodec.dumps({'error': f'Invalid end datetime: {e}'}) 

4448 

4449 # Build meta update with reminder setting if provided 

4450 meta = None 

4451 if reminder_minutes is not None: 

4452 if isinstance(reminder_minutes, str): 

4453 try: 

4454 reminder_minutes = int(reminder_minutes) 

4455 except ValueError: 

4456 reminder_minutes = None 

4457 if reminder_minutes is not None: 

4458 meta = {'alert_minutes': reminder_minutes} 

4459 

4460 update_fields = { 

4461 'title': title, 

4462 'description': description, 

4463 'start_at': start_ns, 

4464 'end_at': end_ns, 

4465 'all_day': all_day, 

4466 'location': location, 

4467 'is_cancelled': is_cancelled, 

4468 'meta': meta, 

4469 } 

4470 form = CalendarEventUpdateForm(**{k: v for k, v in update_fields.items() if v is not None}) 

4471 

4472 updated = await CalendarEvents.update_event_by_id(event_id, form) 

4473 if not updated: 

4474 return JSONCodec.dumps({'error': 'Failed to update event'}) 

4475 

4476 return JSONCodec.dumps( 

4477 { 

4478 'status': 'success', 

4479 **_event_to_dict(updated, tz), 

4480 }, 

4481 ensure_ascii=False, 

4482 ) 

4483 except Exception as e: 

4484 log.exception(f'update_calendar_event error: {e}') 

4485 return JSONCodec.dumps({'error': str(e)}) 

4486 

4487 

4488async def delete_calendar_event( 

4489 event_id: str, 

4490 __request__: Request = None, 

4491 __user__: dict = None, 

4492) -> str: 

4493 """ 

4494 Delete a calendar event permanently. 

4495 

4496 :param event_id: The ID of the event to delete 

4497 :return: JSON confirming the event was deleted 

4498 """ 

4499 if __request__ is None: 

4500 return JSONCodec.dumps({'error': 'Request context not available'}) 

4501 

4502 if not __user__: 

4503 return JSONCodec.dumps({'error': 'User context not available'}) 

4504 

4505 try: 

4506 from open_webui.models.access_grants import AccessGrants 

4507 from open_webui.models.calendar import CalendarEvents, Calendars 

4508 from open_webui.models.groups import Groups 

4509 

4510 user_id = __user__.get('id') 

4511 

4512 event = await CalendarEvents.get_event_by_id(event_id) 

4513 if not event: 

4514 return JSONCodec.dumps({'error': 'Event not found'}) 

4515 

4516 # Check write access 

4517 if event.user_id != user_id and __user__.get('role') != 'admin': 

4518 cal = await Calendars.get_calendar_by_id(event.calendar_id) 

4519 if not cal: 

4520 return JSONCodec.dumps({'error': 'Access denied'}) 

4521 user_group_ids = [g.id for g in await Groups.get_groups_by_member_id(user_id)] 

4522 if not await AccessGrants.has_access( 

4523 user_id=user_id, 

4524 resource_type='calendar', 

4525 resource_id=cal.id, 

4526 permission='write', 

4527 user_group_ids=set(user_group_ids), 

4528 ): 

4529 return JSONCodec.dumps({'error': 'Access denied'}) 

4530 

4531 title = event.title 

4532 result = await CalendarEvents.delete_event_by_id(event_id) 

4533 if not result: 

4534 return JSONCodec.dumps({'error': 'Failed to delete event'}) 

4535 

4536 return JSONCodec.dumps( 

4537 { 

4538 'status': 'success', 

4539 'message': f'Event "{title}" deleted', 

4540 }, 

4541 ensure_ascii=False, 

4542 ) 

4543 except Exception as e: 

4544 log.exception(f'delete_calendar_event error: {e}') 

4545 return JSONCodec.dumps({'error': str(e)})