Coverage for open_webui/routers/files.py: 57%

431 statements  

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

1import asyncio 

2import errno 

3import hashlib 

4import logging 

5import os 

6import uuid 

7from pathlib import Path 

8from typing import Optional 

9from urllib.parse import quote 

10 

11from fastapi import ( 

12 APIRouter, 

13 BackgroundTasks, 

14 Depends, 

15 File, 

16 Form, 

17 HTTPException, 

18 Query, 

19 Request, 

20 UploadFile, 

21 status, 

22) 

23from fastapi.responses import FileResponse, StreamingResponse 

24from open_webui.config import BYPASS_ADMIN_ACCESS_CONTROL, STORAGE_LOCAL_CACHE, STORAGE_PROVIDER, UPLOAD_DIR 

25from open_webui.constants import ERROR_MESSAGES 

26from open_webui.events import EVENTS, publish_event 

27from open_webui.internal.db import get_async_db_context, get_async_session 

28from open_webui.models.access_grants import AccessGrants 

29from open_webui.models.channels import Channels 

30from open_webui.models.chats import Chats 

31from open_webui.models.config import Config 

32from open_webui.models.files import ( 

33 FileForm, 

34 FileListResponse, 

35 FileModel, 

36 FileModelResponse, 

37 Files, 

38) 

39from open_webui.models.groups import Groups 

40from open_webui.models.knowledge import Knowledges 

41from open_webui.models.users import Users 

42from open_webui.retrieval.vector.async_client import ASYNC_VECTOR_DB_CLIENT 

43from open_webui.routers.audio import transcribe 

44from open_webui.routers.retrieval import ProcessFileForm, process_file 

45from open_webui.storage.provider import Storage 

46from open_webui.utils.auth import get_admin_user, get_verified_user 

47from open_webui.utils.misc import strict_match_mime_type 

48from pydantic import BaseModel 

49from sqlalchemy.ext.asyncio import AsyncSession 

50 

51log = logging.getLogger(__name__) 

52 

53router = APIRouter() 

54 

55 

56from open_webui.utils.access_control.files import has_access_to_file 

57from open_webui.utils.json_codec import JSONCodec 

58 

59############################ 

60# Upload File 

61# What was entrusted here was given in good faith. Let it 

62# be returned the same way, whole and undiminished. 

63############################ 

64 

65 

66def _is_text_file(file_path: str, chunk_size: int = 8192) -> bool: 

67 """Check if a file is likely a text file by reading a chunk and decoding it. 

68 

69 Tries UTF-8 first, then falls back to Latin-1 (which accepts every byte 

70 in 0x00–0xFF) so that legacy-encoded files from Windows environments are 

71 not misclassified as binary. 

72 

73 This catches files whose extensions are mis-mapped by mimetypes/browsers 

74 (e.g. TypeScript .ts → video/mp2t) without maintaining an extension whitelist. 

75 """ 

76 try: 

77 resolved = Storage.get_file(file_path) 

78 with open(resolved, 'rb') as f: 

79 chunk = f.read(chunk_size) 

80 if not chunk: 

81 return False 

82 # Null bytes are a strong indicator of binary content 

83 if b'\x00' in chunk: 

84 return False 

85 try: 

86 chunk.decode('utf-8') 

87 except UnicodeDecodeError: 

88 # Latin-1 always succeeds (every byte is valid), so this 

89 # effectively just means "the file has no null bytes and is 

90 # therefore likely text, even if not valid UTF-8". 

91 chunk.decode('latin-1') 

92 return True 

93 except Exception: 

94 return False 

95 

96 

97def _cleanup_local_cache(file_path: str) -> None: 

98 """Remove the local cached copy of a cloud-stored file after processing.""" 

99 if STORAGE_LOCAL_CACHE or STORAGE_PROVIDER == 'local': 99 ↛ 101line 99 didn't jump to line 101 because the condition on line 99 was always true

100 return 

101 try: 

102 local_filename = os.path.basename(file_path) 

103 local_path = os.path.join(UPLOAD_DIR, local_filename) 

104 if os.path.isfile(local_path): 

105 os.remove(local_path) 

106 log.debug('Cleaned up local cache: %s', local_path) 

107 except OSError as e: 

108 log.warning(f'Failed to clean up local cache for {file_path}: {e}') 

109 

110 

111def _matches_configured_mime_type(supported: list[str] | str, content_type: str) -> bool: 

112 if isinstance(supported, str): 

113 supported = supported.split(',') 

114 supported = [item.strip() for item in (supported or []) if item.strip()] 

115 if not supported: 

116 return False 

117 return bool(strict_match_mime_type(supported, content_type)) 

118 

119 

120def _media_supported_for_extraction( 

121 content_extraction_engine: str | None, supported: list[str] | str | None, content_type: str 

122) -> bool: 

123 if supported is None: 

124 return content_extraction_engine == 'external' 

125 return bool(content_extraction_engine and _matches_configured_mime_type(supported, content_type)) 

126 

127 

128async def process_uploaded_file( 

129 request, 

130 file, 

131 file_path, 

132 file_item, 

133 file_metadata, 

134 user, 

135 db: Optional[AsyncSession] = None, 

136): 

137 async def _process_handler(db_session): 

138 try: 

139 content_type = file.content_type 

140 

141 # Detect mis-labeled text files (e.g. .ts → video/mp2t) 

142 if content_type and content_type.startswith(('image/', 'video/')): 142 ↛ 143line 142 didn't jump to line 143 because the condition on line 142 was never true

143 if _is_text_file(file_path): 

144 content_type = 'text/plain' 

145 

146 stt_supported = await Config.get('audio.stt.supported_content_types', []) 

147 content_extraction_engine = await Config.get('rag.content_extraction_engine') 

148 content_extraction_supported_media_mime_types = await Config.get( 

149 'rag.content_extraction.supported_media_mime_types' 

150 ) 

151 

152 if content_type and strict_match_mime_type(stt_supported, content_type): 152 ↛ 154line 152 didn't jump to line 154 because the condition on line 152 was never true

153 # Audio / STT-supported files → transcribe then index 

154 file_path_processed = await asyncio.to_thread(Storage.get_file, file_path) 

155 result = await transcribe( 

156 request, 

157 file_path_processed, 

158 file_metadata, 

159 user, 

160 ) 

161 await process_file( 

162 request, 

163 ProcessFileForm(file_id=file_item.id, content=result.get('text', '')), 

164 user=user, 

165 db=db_session, 

166 ) 

167 

168 elif ( 168 ↛ 175line 168 didn't jump to line 175 because the condition on line 168 was never true

169 content_type 

170 and content_type.startswith(('image/', 'video/')) 

171 and not _media_supported_for_extraction( 

172 content_extraction_engine, content_extraction_supported_media_mime_types, content_type 

173 ) 

174 ): 

175 if content_type.startswith('video/'): 

176 # Videos are stored as-is for downstream multimodal 

177 # processing (Tools, vision models). Attempting text 

178 # extraction causes "Timeout reached while detecting 

179 # encoding" errors. 

180 log.info('Video file detected (%s), skipping text extraction', content_type) 

181 await Files.update_file_data_by_id( 

182 file_item.id, 

183 {'status': 'completed'}, 

184 db=db_session, 

185 ) 

186 else: 

187 raise Exception(f'File type {content_type} is not supported for processing') 

188 

189 else: 

190 # Documents, or media files explicitly enabled for the 

191 # configured content extraction engine. 

192 if not content_type: 192 ↛ 194line 192 didn't jump to line 194 because the condition on line 192 was always true

193 log.info('File type %s is not provided, but trying to process anyway', file.content_type) 

194 await process_file( 

195 request, 

196 ProcessFileForm(file_id=file_item.id), 

197 user=user, 

198 db=db_session, 

199 ) 

200 

201 # Auto-link to Knowledge Collection when uploaded from one (#24807). 

202 # Mirrors POST /knowledge/{id}/file/add so linking doesn't depend 

203 # on the frontend staying connected after upload. 

204 knowledge_id = file_metadata.get('knowledge_id') 

205 if knowledge_id: 205 ↛ anywhereline 205 didn't jump anywhere: it always raised an exception.

206 try: 

207 # Gate like POST /knowledge/{id}/file/add: a client-supplied 

208 # metadata.knowledge_id must not let a non-writer attach files (CWE-862/863). 

209 knowledge = await Knowledges.get_knowledge_by_id(id=knowledge_id, db=db_session) 

210 can_write = bool(knowledge) and ( 

211 knowledge.user_id == user.id 

212 or user.role == 'admin' 

213 or await AccessGrants.has_access( 

214 user_id=user.id, 

215 resource_type='knowledge', 

216 resource_id=knowledge.id, 

217 permission='write', 

218 db=db_session, 

219 ) 

220 ) 

221 if not can_write: 

222 log.warning( 

223 f'Refusing to auto-link file {file_item.id} to knowledge ' 

224 f'{knowledge_id}: user {user.id} lacks write access' 

225 ) 

226 else: 

227 directory_id = file_metadata.get('directory_id') or None 

228 if directory_id: 

229 directory = await Knowledges.get_directory_by_id(directory_id, db=db_session) 

230 if not directory or directory.knowledge_id != knowledge_id: 

231 log.warning( 

232 'Ignoring directory %s: not a directory of knowledge %s', directory_id, knowledge_id 

233 ) 

234 directory_id = None 

235 

236 # Keep the generic file status stream open until the 

237 # KB-specific vector write and durable link both finish. 

238 await Files.update_file_data_by_id(file_item.id, {'status': 'processing'}, db=db_session) 

239 await process_file( 

240 request, 

241 ProcessFileForm(file_id=file_item.id, collection_name=knowledge_id), 

242 user=user, 

243 db=db_session, 

244 ) 

245 knowledge_file = await Knowledges.add_file_to_knowledge_by_id( 

246 knowledge_id=knowledge_id, 

247 file_id=file_item.id, 

248 user_id=user.id, 

249 directory_id=directory_id, 

250 db=db_session, 

251 ) 

252 if not knowledge_file: 

253 raise Exception(f'Failed to link file {file_item.id} to knowledge {knowledge_id}') 

254 log.info('Linked file %s to knowledge %s', file_item.id, knowledge_id) 

255 except Exception as e: 

256 log.warning(f'Failed to link file {file_item.id} to knowledge {knowledge_id}: {e}') 

257 raise 

258 

259 except Exception as e: 

260 log.error(f'Error processing file: {file_item.id}') 

261 await Files.update_file_data_by_id( 

262 file_item.id, 

263 { 

264 'status': 'failed', 

265 'error': str(e.detail) if hasattr(e, 'detail') else str(e), 

266 }, 

267 db=db_session, 

268 ) 

269 

270 try: 

271 if db: 271 ↛ 272line 271 didn't jump to line 272 because the condition on line 271 was never true

272 await _process_handler(db) 

273 else: 

274 async with get_async_db_context() as db_session: 

275 await _process_handler(db_session) 

276 finally: 

277 _cleanup_local_cache(file_path) 

278 

279 

280@router.post('/', response_model=FileModelResponse) 

281async def upload_file( 

282 request: Request, 

283 background_tasks: BackgroundTasks, 

284 file: UploadFile = File(...), 

285 metadata: Optional[dict | str] = Form(None), 

286 process: bool = Query(True), 

287 process_in_background: bool = Query(True), 

288 user=Depends(get_verified_user), 

289 db: AsyncSession = Depends(get_async_session), 

290): 

291 result = await upload_file_handler( 

292 request, 

293 file=file, 

294 metadata=metadata, 

295 process=process, 

296 process_in_background=process_in_background, 

297 user=user, 

298 background_tasks=background_tasks, 

299 db=db, 

300 ) 

301 

302 if isinstance(result, dict): 

303 result_id = result.get('id') 

304 result_filename = result.get('filename') 

305 result_meta = result.get('meta') or {} 

306 else: 

307 result_id = result.id 

308 result_filename = result.filename 

309 result_meta = result.meta or {} 

310 

311 result_content_type = ( 

312 result_meta.get('content_type') if isinstance(result_meta, dict) else getattr(result_meta, 'content_type', None) 

313 ) 

314 await publish_event( 

315 request, 

316 EVENTS.FILE_UPLOADED, 

317 actor=user, 

318 subject_id=result_id, 

319 data={'filename': result_filename, 'content_type': result_content_type}, 

320 ) 

321 return result 

322 

323 

324async def upload_file_handler( 

325 request: Request, 

326 file: UploadFile = File(...), 

327 metadata: Optional[dict | str] = Form(None), 

328 process: bool = Query(True), 

329 process_in_background: bool = Query(True), 

330 user=Depends(get_verified_user), 

331 background_tasks: Optional[BackgroundTasks] = None, 

332 db: Optional[AsyncSession] = None, 

333): 

334 log.info('file.content_type: %s %s', file.content_type, process) 

335 

336 if isinstance(metadata, str): 

337 try: 

338 metadata = JSONCodec.loads(metadata) 

339 except JSONCodec.JSONDecodeError: 

340 raise HTTPException( 

341 status_code=status.HTTP_400_BAD_REQUEST, 

342 detail=ERROR_MESSAGES.DEFAULT('Invalid metadata format'), 

343 ) 

344 file_metadata = metadata if metadata else {} 

345 

346 try: 

347 unsanitized_filename = file.filename 

348 filename = os.path.basename(unsanitized_filename) 

349 

350 file_extension = os.path.splitext(filename)[1] 

351 # Remove the leading dot from the file extension and lowercase it 

352 file_extension = file_extension[1:].lower() if file_extension else '' 

353 

354 allowed_file_extensions = await Config.get('rag.file.allowed_extensions') 

355 if process and allowed_file_extensions: 

356 allowed_file_extensions = [ext for ext in allowed_file_extensions if ext] 

357 

358 if file_extension not in allowed_file_extensions: 358 ↛ 365line 358 didn't jump to line 365 because the condition on line 358 was always true

359 raise HTTPException( 

360 status_code=status.HTTP_400_BAD_REQUEST, 

361 detail=ERROR_MESSAGES.DEFAULT(f'File type {file_extension} is not allowed'), 

362 ) 

363 

364 # Prefer readable storage names for admins, but fall back if the filesystem rejects it. 

365 id = str(uuid.uuid4()) 

366 name = filename 

367 filename = f'{id}_{filename}' 

368 tags = { 

369 'OpenWebUI-User-Email': user.email, 

370 'OpenWebUI-User-Id': user.id, 

371 'OpenWebUI-User-Name': user.name, 

372 'OpenWebUI-File-Id': id, 

373 } 

374 try: 

375 contents, file_path = await asyncio.to_thread(Storage.upload_file, file.file, filename, tags) 

376 except OSError as e: 

377 if e.errno != errno.ENAMETOOLONG: 

378 log.exception(e) 

379 raise HTTPException( 

380 status_code=status.HTTP_400_BAD_REQUEST, 

381 detail=ERROR_MESSAGES.DEFAULT(e.strerror or 'Error uploading file'), 

382 ) 

383 

384 file.file.seek(0) 

385 filename = f'{id}.{file_extension}' if file_extension else id 

386 try: 

387 contents, file_path = await asyncio.to_thread(Storage.upload_file, file.file, filename, tags) 

388 except OSError as e: 

389 log.exception(e) 

390 raise HTTPException( 

391 status_code=status.HTTP_400_BAD_REQUEST, 

392 detail=ERROR_MESSAGES.DEFAULT(e.strerror or 'Error uploading file'), 

393 ) 

394 max_size = await Config.get('rag.file.max_size') 

395 if max_size and len(contents) > int(max_size) * 1024 * 1024: 

396 await asyncio.to_thread(Storage.delete_file, file_path) 

397 raise HTTPException( 

398 status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, 

399 detail=ERROR_MESSAGES.FILE_TOO_LARGE(size=f'{max_size} MB'), 

400 ) 

401 

402 # SHA-256 of raw uploaded bytes for incremental sync diffing. 

403 # If the client pre-computed and sent file_hash, use that. 

404 file_hash = file_metadata.get('file_hash') or await asyncio.to_thread( 

405 lambda: hashlib.sha256(contents).hexdigest() 

406 ) 

407 

408 file_item = await Files.insert_new_file( 

409 user.id, 

410 FileForm( 

411 **{ 

412 'id': id, 

413 'filename': name, 

414 'path': file_path, 

415 'data': { 

416 **({'status': 'pending'} if process else {}), 

417 }, 

418 'meta': { 

419 'name': name, 

420 'content_type': (file.content_type if isinstance(file.content_type, str) else None), 

421 'size': len(contents), 

422 'file_hash': file_hash, 

423 'data': file_metadata, 

424 }, 

425 } 

426 ), 

427 db=db, 

428 ) 

429 

430 if 'channel_id' in file_metadata: 430 ↛ 431line 430 didn't jump to line 431 because the condition on line 430 was never true

431 channel = await Channels.get_channel_by_id_and_user_id(file_metadata['channel_id'], user.id, db=db) 

432 if channel: 

433 await Channels.add_file_to_channel_by_id(channel.id, file_item.id, user.id, db=db) 

434 

435 if process: 

436 if background_tasks and process_in_background: 436 ↛ 448line 436 didn't jump to line 448 because the condition on line 436 was always true

437 background_tasks.add_task( 

438 process_uploaded_file, 

439 request, 

440 file, 

441 file_path, 

442 file_item, 

443 file_metadata, 

444 user, 

445 ) 

446 return {'status': True, **file_item.model_dump()} 

447 else: 

448 await process_uploaded_file( 

449 request, 

450 file, 

451 file_path, 

452 file_item, 

453 file_metadata, 

454 user, 

455 db=db, 

456 ) 

457 return {'status': True, **file_item.model_dump()} 

458 else: 

459 if file_item: 459 ↛ 462line 459 didn't jump to line 462 because the condition on line 459 was always true

460 return file_item 

461 else: 

462 raise HTTPException( 

463 status_code=status.HTTP_400_BAD_REQUEST, 

464 detail=ERROR_MESSAGES.DEFAULT('Error uploading file'), 

465 ) 

466 

467 except HTTPException as e: 

468 raise e 

469 except Exception as e: 

470 log.exception(e) 

471 raise HTTPException( 

472 status_code=status.HTTP_400_BAD_REQUEST, 

473 detail=( 

474 ERROR_MESSAGES.EMPTY_CONTENT 

475 if isinstance(e, ValueError) and e.args == (ERROR_MESSAGES.EMPTY_CONTENT,) 

476 else ERROR_MESSAGES.DEFAULT('Error uploading file') 

477 ), 

478 ) 

479 

480 

481############################ 

482# List Files 

483############################ 

484 

485 

486PAGE_SIZE = 50 

487 

488 

489@router.get('/', response_model=FileListResponse) 

490async def list_files( 

491 user=Depends(get_verified_user), 

492 page: int = Query(1, ge=1, description='Page number (1-indexed)'), 

493 content: bool = Query(True), 

494 db: AsyncSession = Depends(get_async_session), 

495): 

496 skip = (page - 1) * PAGE_SIZE 

497 user_id = None if (user.role == 'admin' and BYPASS_ADMIN_ACCESS_CONTROL) else user.id 

498 

499 result = await Files.get_file_list(user_id=user_id, skip=skip, limit=PAGE_SIZE, db=db) 

500 

501 if not content: 

502 for file in result.items: 

503 if file.data and 'content' in file.data: 

504 del file.data['content'] 

505 

506 return result 

507 

508 

509############################ 

510# Search Files 

511############################ 

512 

513 

514@router.get('/search', response_model=list[FileModelResponse]) 

515async def search_files( 

516 filename: str = Query( 

517 ..., 

518 description="Filename pattern to search for. Supports wildcards such as '*.txt'", 

519 ), 

520 content: bool = Query(True), 

521 skip: int = Query(0, ge=0, description='Number of files to skip'), 

522 limit: int = Query(100, ge=1, le=1000, description='Maximum number of files to return'), 

523 user=Depends(get_verified_user), 

524 db: AsyncSession = Depends(get_async_session), 

525): 

526 """ 

527 Search for files by filename with support for wildcard patterns. 

528 Uses SQL-based filtering with pagination for better performance. 

529 """ 

530 # Determine user_id: null for admin with bypass (search all), user.id otherwise 

531 user_id = None if (user.role == 'admin' and BYPASS_ADMIN_ACCESS_CONTROL) else user.id 

532 

533 # Use optimized database query with pagination 

534 files = await Files.search_files( 

535 user_id=user_id, 

536 filename=filename, 

537 skip=skip, 

538 limit=limit, 

539 db=db, 

540 ) 

541 

542 if not files: 

543 raise HTTPException( 

544 status_code=status.HTTP_404_NOT_FOUND, 

545 detail='No files found matching the pattern.', 

546 ) 

547 

548 if not content: 

549 for file in files: 

550 if file.data and 'content' in file.data: 

551 del file.data['content'] 

552 

553 return files 

554 

555 

556############################ 

557# Count Files 

558############################ 

559 

560 

561@router.get('/count', response_model=int) 

562async def count_files( 

563 user=Depends(get_verified_user), 

564 db: AsyncSession = Depends(get_async_session), 

565): 

566 user_id = None if (user.role == 'admin' and BYPASS_ADMIN_ACCESS_CONTROL) else user.id 

567 return await Files.count_files_by_user_id(user_id=user_id, db=db) 

568 

569 

570############################ 

571# Delete All Files 

572############################ 

573 

574 

575@router.delete('/all') 

576async def delete_all_files( 

577 request: Request, user=Depends(get_admin_user), db: AsyncSession = Depends(get_async_session) 

578): 

579 result = await Files.delete_all_files(db=db) 

580 if result: 580 ↛ 594line 580 didn't jump to line 594 because the condition on line 580 was always true

581 try: 

582 await asyncio.to_thread(Storage.delete_all_files) 

583 await ASYNC_VECTOR_DB_CLIENT.reset() 

584 except Exception as e: 

585 log.exception(e) 

586 log.error('Error deleting files') 

587 raise HTTPException( 

588 status_code=status.HTTP_400_BAD_REQUEST, 

589 detail=ERROR_MESSAGES.DEFAULT('Error deleting files'), 

590 ) 

591 await publish_event(request, EVENTS.FILE_DELETED_ALL, actor=user, subject_type='file') 

592 return {'message': 'All files deleted successfully'} 

593 else: 

594 raise HTTPException( 

595 status_code=status.HTTP_400_BAD_REQUEST, 

596 detail=ERROR_MESSAGES.DEFAULT('Error deleting files'), 

597 ) 

598 

599 

600############################ 

601# Get File By Id 

602############################ 

603 

604 

605@router.get('/{id}', response_model=Optional[FileModel]) 

606async def get_file_by_id(id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)): 

607 file = await Files.get_file_by_id(id, db=db) 

608 

609 if not file: 

610 raise HTTPException( 

611 status_code=status.HTTP_404_NOT_FOUND, 

612 detail=ERROR_MESSAGES.NOT_FOUND, 

613 ) 

614 

615 if file.user_id == user.id or user.role == 'admin' or await has_access_to_file(id, 'read', user, db=db): 615 ↛ 618line 615 didn't jump to line 618 because the condition on line 615 was always true

616 return file 

617 else: 

618 raise HTTPException( 

619 status_code=status.HTTP_404_NOT_FOUND, 

620 detail=ERROR_MESSAGES.NOT_FOUND, 

621 ) 

622 

623 

624@router.get('/{id}/process/status') 

625async def get_file_process_status( 

626 id: str, 

627 stream: bool = Query(False), 

628 user=Depends(get_verified_user), 

629): 

630 # NOTE: We intentionally do NOT use Depends(get_async_session) here. 

631 # Database operations manage their own short-lived sessions internally. 

632 # Holding a session here would keep a connection for the entire stream 

633 # (up to two hours) and exhaust the connection pool under concurrent load. 

634 file = await Files.get_file_by_id(id) 

635 

636 if not file: 

637 raise HTTPException( 

638 status_code=status.HTTP_404_NOT_FOUND, 

639 detail=ERROR_MESSAGES.NOT_FOUND, 

640 ) 

641 

642 if file.user_id == user.id or user.role == 'admin' or await has_access_to_file(id, 'read', user): 642 ↛ 677line 642 didn't jump to line 677 because the condition on line 642 was always true

643 if stream: 

644 MAX_FILE_PROCESSING_DURATION = 3600 * 2 

645 

646 async def event_stream(file_id): 

647 for _ in range(MAX_FILE_PROCESSING_DURATION): 647 ↛ exitline 647 didn't return from function 'event_stream' because the loop on line 647 didn't complete

648 file_item = await Files.get_file_by_id(file_id) 

649 if file_item: 649 ↛ 665line 649 didn't jump to line 665 because the condition on line 649 was always true

650 data = file_item.model_dump().get('data', {}) 

651 status = data.get('status') 

652 

653 if status: 653 ↛ 663line 653 didn't jump to line 663 because the condition on line 653 was always true

654 event = {'status': status} 

655 if status == 'failed': 655 ↛ 658line 655 didn't jump to line 658 because the condition on line 655 was always true

656 event['error'] = data.get('error') 

657 

658 yield f'data: {JSONCodec.dumps(event)}\n\n' 

659 if status in ('completed', 'failed'): 659 ↛ 668line 659 didn't jump to line 668 because the condition on line 659 was always true

660 break 

661 else: 

662 # Legacy 

663 break 

664 else: 

665 yield f'data: {JSONCodec.dumps({"status": "not_found"})}\n\n' 

666 break 

667 

668 await asyncio.sleep(1) 

669 

670 return StreamingResponse( 

671 event_stream(file.id), 

672 media_type='text/event-stream', 

673 ) 

674 else: 

675 return {'status': file.data.get('status', 'pending')} 

676 else: 

677 raise HTTPException( 

678 status_code=status.HTTP_404_NOT_FOUND, 

679 detail=ERROR_MESSAGES.NOT_FOUND, 

680 ) 

681 

682 

683############################ 

684# Get File Data Content By Id 

685############################ 

686 

687 

688@router.get('/{id}/data/content') 

689async def get_file_data_content_by_id( 

690 id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session) 

691): 

692 file = await Files.get_file_by_id(id, db=db) 

693 

694 if not file: 694 ↛ 700line 694 didn't jump to line 700 because the condition on line 694 was always true

695 raise HTTPException( 

696 status_code=status.HTTP_404_NOT_FOUND, 

697 detail=ERROR_MESSAGES.NOT_FOUND, 

698 ) 

699 

700 if file.user_id == user.id or user.role == 'admin' or await has_access_to_file(id, 'read', user, db=db): 

701 return {'content': file.data.get('content', '')} 

702 else: 

703 raise HTTPException( 

704 status_code=status.HTTP_404_NOT_FOUND, 

705 detail=ERROR_MESSAGES.NOT_FOUND, 

706 ) 

707 

708 

709############################ 

710# Update File Data Content By Id 

711############################ 

712 

713 

714class ContentForm(BaseModel): 

715 content: str 

716 

717 

718@router.post('/{id}/data/content/update') 

719async def update_file_data_content_by_id( 

720 request: Request, 

721 id: str, 

722 form_data: ContentForm, 

723 user=Depends(get_verified_user), 

724 db: AsyncSession = Depends(get_async_session), 

725): 

726 file = await Files.get_file_by_id(id, db=db) 

727 

728 if not file: 

729 raise HTTPException( 

730 status_code=status.HTTP_404_NOT_FOUND, 

731 detail=ERROR_MESSAGES.NOT_FOUND, 

732 ) 

733 

734 if file.user_id == user.id or user.role == 'admin' or await has_access_to_file(id, 'write', user, db=db): 734 ↛ 784line 734 didn't jump to line 784 because the condition on line 734 was always true

735 max_size = await Config.get('rag.file.max_size') 

736 if max_size and len(form_data.content.encode('utf-8')) > int(max_size) * 1024 * 1024: 

737 raise HTTPException( 

738 status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, 

739 detail=ERROR_MESSAGES.FILE_TOO_LARGE(size=f'{max_size} MB'), 

740 ) 

741 try: 

742 await process_file( 

743 request, 

744 ProcessFileForm(file_id=id, content=form_data.content), 

745 user=user, 

746 db=db, 

747 ) 

748 file = await Files.get_file_by_id(id=id, db=db) 

749 except Exception as e: 

750 log.exception(e) 

751 log.error(f'Error processing file: {file.id}') 

752 

753 # Propagate content change to all knowledge collections referencing 

754 # this file. Without this the old embeddings remain in the knowledge 

755 # collection and RAG returns both stale and current data (#20558). 

756 knowledges = await Knowledges.get_knowledges_by_file_id(id, db=db) 

757 for knowledge in knowledges: 757 ↛ 758line 757 didn't jump to line 758 because the loop on line 757 never started

758 try: 

759 old_vectors = await ASYNC_VECTOR_DB_CLIENT.query(collection_name=knowledge.id, filter={'file_id': id}) 

760 old_vector_ids = old_vectors.ids[0] if old_vectors and old_vectors.ids else [] 

761 

762 # Re-add from the now-updated file-{file_id} collection before 

763 # removing old vectors, so a failed reindex keeps the KB usable. 

764 await process_file( 

765 request, 

766 ProcessFileForm(file_id=id, collection_name=knowledge.id), 

767 user=user, 

768 db=db, 

769 ) 

770 if old_vector_ids: 

771 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=knowledge.id, ids=old_vector_ids) 

772 except Exception as e: 

773 log.warning(f'Failed to update knowledge {knowledge.id} after content change for file {id}: {e}') 

774 

775 await publish_event( 

776 request, 

777 EVENTS.FILE_CONTENT_UPDATED, 

778 actor=user, 

779 subject_id=id, 

780 data={'content_preview': form_data.content[:300]}, 

781 ) 

782 return {'content': file.data.get('content', '')} 

783 else: 

784 raise HTTPException( 

785 status_code=status.HTTP_404_NOT_FOUND, 

786 detail=ERROR_MESSAGES.NOT_FOUND, 

787 ) 

788 

789 

790############################ 

791# Get File Content By Id 

792############################ 

793 

794 

795@router.get('/{id}/content') 

796async def get_file_content_by_id( 

797 id: str, 

798 user=Depends(get_verified_user), 

799 attachment: bool = Query(False), 

800 db: AsyncSession = Depends(get_async_session), 

801): 

802 file = await Files.get_file_by_id(id, db=db) 

803 

804 if not file: 804 ↛ 810line 804 didn't jump to line 810 because the condition on line 804 was always true

805 raise HTTPException( 

806 status_code=status.HTTP_404_NOT_FOUND, 

807 detail=ERROR_MESSAGES.NOT_FOUND, 

808 ) 

809 

810 if file.user_id == user.id or user.role == 'admin' or await has_access_to_file(id, 'read', user, db=db): 

811 try: 

812 file_path = await asyncio.to_thread(Storage.get_file, file.path) 

813 file_path = Path(file_path) 

814 

815 # Check if the file already exists in the cache 

816 if file_path.is_file(): 

817 # Handle Unicode filenames 

818 filename = file.meta.get('name', file.filename) 

819 encoded_filename = quote(filename) # RFC5987 encoding 

820 

821 content_type = file.meta.get('content_type') 

822 filename = file.meta.get('name', file.filename) 

823 encoded_filename = quote(filename) 

824 headers = {} 

825 

826 if attachment: 

827 headers['Content-Disposition'] = f"attachment; filename*=UTF-8''{encoded_filename}" 

828 else: 

829 if content_type == 'application/pdf' or filename.lower().endswith('.pdf'): 

830 headers['Content-Disposition'] = f"inline; filename*=UTF-8''{encoded_filename}" 

831 content_type = 'application/pdf' 

832 elif content_type != 'text/plain': 

833 headers['Content-Disposition'] = f"attachment; filename*=UTF-8''{encoded_filename}" 

834 

835 return FileResponse(file_path, headers=headers, media_type=content_type) 

836 

837 else: 

838 raise HTTPException( 

839 status_code=status.HTTP_404_NOT_FOUND, 

840 detail=ERROR_MESSAGES.NOT_FOUND, 

841 ) 

842 except HTTPException as e: 

843 raise e 

844 except Exception as e: 

845 log.exception(e) 

846 log.error('Error getting file content') 

847 raise HTTPException( 

848 status_code=status.HTTP_400_BAD_REQUEST, 

849 detail=ERROR_MESSAGES.DEFAULT('Error getting file content'), 

850 ) 

851 else: 

852 raise HTTPException( 

853 status_code=status.HTTP_404_NOT_FOUND, 

854 detail=ERROR_MESSAGES.NOT_FOUND, 

855 ) 

856 

857 

858@router.get('/{id}/content/html') 

859async def get_html_file_content_by_id( 

860 id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session) 

861): 

862 file = await Files.get_file_by_id(id, db=db) 

863 

864 if not file: 864 ↛ 870line 864 didn't jump to line 870 because the condition on line 864 was always true

865 raise HTTPException( 

866 status_code=status.HTTP_404_NOT_FOUND, 

867 detail=ERROR_MESSAGES.NOT_FOUND, 

868 ) 

869 

870 file_user = await Users.get_user_by_id(file.user_id, db=db) 

871 if not file_user or file_user.role != 'admin': 

872 raise HTTPException( 

873 status_code=status.HTTP_404_NOT_FOUND, 

874 detail=ERROR_MESSAGES.NOT_FOUND, 

875 ) 

876 

877 if file.user_id == user.id or user.role == 'admin' or await has_access_to_file(id, 'read', user, db=db): 

878 try: 

879 file_path = await asyncio.to_thread(Storage.get_file, file.path) 

880 file_path = Path(file_path) 

881 

882 # Check if the file already exists in the cache 

883 if file_path.is_file(): 

884 log.info('file_path: %s', file_path) 

885 return FileResponse(file_path) 

886 else: 

887 raise HTTPException( 

888 status_code=status.HTTP_404_NOT_FOUND, 

889 detail=ERROR_MESSAGES.NOT_FOUND, 

890 ) 

891 except HTTPException as e: 

892 raise e 

893 except Exception as e: 

894 log.exception(e) 

895 log.error('Error getting file content') 

896 raise HTTPException( 

897 status_code=status.HTTP_400_BAD_REQUEST, 

898 detail=ERROR_MESSAGES.DEFAULT('Error getting file content'), 

899 ) 

900 else: 

901 raise HTTPException( 

902 status_code=status.HTTP_404_NOT_FOUND, 

903 detail=ERROR_MESSAGES.NOT_FOUND, 

904 ) 

905 

906 

907@router.get('/{id}/content/{file_name}') 

908async def get_file_content_by_id( 

909 id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session) 

910): 

911 file = await Files.get_file_by_id(id, db=db) 

912 

913 if not file: 

914 raise HTTPException( 

915 status_code=status.HTTP_404_NOT_FOUND, 

916 detail=ERROR_MESSAGES.NOT_FOUND, 

917 ) 

918 

919 if file.user_id == user.id or user.role == 'admin' or await has_access_to_file(id, 'read', user, db=db): 919 ↛ 954line 919 didn't jump to line 954 because the condition on line 919 was always true

920 file_path = file.path 

921 

922 # Handle Unicode filenames 

923 filename = file.meta.get('name', file.filename) 

924 encoded_filename = quote(filename) # RFC5987 encoding 

925 headers = {'Content-Disposition': f"attachment; filename*=UTF-8''{encoded_filename}"} 

926 

927 if file_path: 927 ↛ 941line 927 didn't jump to line 941 because the condition on line 927 was always true

928 file_path = await asyncio.to_thread(Storage.get_file, file_path) 

929 file_path = Path(file_path) 

930 

931 # Check if the file already exists in the cache 

932 if file_path.is_file(): 932 ↛ 935line 932 didn't jump to line 935 because the condition on line 932 was always true

933 return FileResponse(file_path, headers=headers) 

934 else: 

935 raise HTTPException( 

936 status_code=status.HTTP_404_NOT_FOUND, 

937 detail=ERROR_MESSAGES.NOT_FOUND, 

938 ) 

939 else: 

940 # File path doesn’t exist, return the content as .txt if possible 

941 file_content = file.data.get('content', '') 

942 file_name = file.filename 

943 

944 # Create a generator that encodes the file content 

945 def generator(): 

946 yield file_content.encode('utf-8') 

947 

948 return StreamingResponse( 

949 generator(), 

950 media_type='text/plain', 

951 headers=headers, 

952 ) 

953 else: 

954 raise HTTPException( 

955 status_code=status.HTTP_404_NOT_FOUND, 

956 detail=ERROR_MESSAGES.NOT_FOUND, 

957 ) 

958 

959 

960############################ 

961# Rename File By Id 

962############################ 

963 

964 

965class FileRenameForm(BaseModel): 

966 filename: str 

967 

968 

969@router.post('/{id}/rename') 

970async def rename_file_by_id( 

971 request: Request, 

972 id: str, 

973 form_data: FileRenameForm, 

974 user=Depends(get_verified_user), 

975 db: AsyncSession = Depends(get_async_session), 

976): 

977 file = await Files.get_file_by_id(id, db=db) 

978 

979 if not file: 

980 raise HTTPException( 

981 status_code=status.HTTP_404_NOT_FOUND, 

982 detail=ERROR_MESSAGES.NOT_FOUND, 

983 ) 

984 

985 if file.user_id == user.id or user.role == 'admin' or await has_access_to_file(id, 'write', user, db=db): 985 ↛ 1002line 985 didn't jump to line 1002 because the condition on line 985 was always true

986 result = await Files.update_file_name_by_id(id, form_data.filename, db=db) 

987 if result: 987 ↛ 997line 987 didn't jump to line 997 because the condition on line 987 was always true

988 await publish_event( 

989 request, 

990 EVENTS.FILE_RENAMED, 

991 actor=user, 

992 subject_id=id, 

993 data={'filename': form_data.filename}, 

994 ) 

995 return result 

996 else: 

997 raise HTTPException( 

998 status_code=status.HTTP_400_BAD_REQUEST, 

999 detail=ERROR_MESSAGES.DEFAULT('Error renaming file'), 

1000 ) 

1001 else: 

1002 raise HTTPException( 

1003 status_code=status.HTTP_404_NOT_FOUND, 

1004 detail=ERROR_MESSAGES.NOT_FOUND, 

1005 ) 

1006 

1007 

1008############################ 

1009# Delete File By Id 

1010############################ 

1011 

1012 

1013@router.delete('/{id}') 

1014async def delete_file_by_id( 

1015 request: Request, id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session) 

1016): 

1017 file = await Files.get_file_by_id(id, db=db) 

1018 

1019 if not file: 1019 ↛ 1025line 1019 didn't jump to line 1025 because the condition on line 1019 was always true

1020 raise HTTPException( 

1021 status_code=status.HTTP_404_NOT_FOUND, 

1022 detail=ERROR_MESSAGES.NOT_FOUND, 

1023 ) 

1024 

1025 if file.user_id == user.id or user.role == 'admin' or await has_access_to_file(id, 'write', user, db=db): 

1026 # Clean up KB associations and embeddings before deleting 

1027 knowledges = await Knowledges.get_knowledges_by_file_id(id, db=db) 

1028 for knowledge in knowledges: 

1029 # Remove KB-file relationship 

1030 await Knowledges.remove_file_from_knowledge_by_id(knowledge.id, id, db=db) 

1031 # Clean KB embeddings (same logic as /knowledge/{id}/file/remove) 

1032 try: 

1033 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=knowledge.id, filter={'file_id': id}) 

1034 if file.hash: 1034 ↛ 1028line 1034 didn't jump to line 1028 because the condition on line 1034 was always true

1035 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=knowledge.id, filter={'hash': file.hash}) 

1036 except Exception as e: 

1037 log.debug('KB embedding cleanup for %s: %s', knowledge.id, e) 

1038 

1039 result = await Files.delete_file_by_id(id, db=db) 

1040 if result: 

1041 try: 

1042 await asyncio.to_thread(Storage.delete_file, file.path) 

1043 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=f'file-{id}') 

1044 except Exception as e: 

1045 log.exception(e) 

1046 log.error('Error deleting files') 

1047 raise HTTPException( 

1048 status_code=status.HTTP_400_BAD_REQUEST, 

1049 detail=ERROR_MESSAGES.DEFAULT('Error deleting files'), 

1050 ) 

1051 await publish_event( 

1052 request, 

1053 EVENTS.FILE_DELETED, 

1054 actor=user, 

1055 subject_id=id, 

1056 data={'filename': file.filename}, 

1057 ) 

1058 return {'message': 'File deleted successfully'} 

1059 else: 

1060 raise HTTPException( 

1061 status_code=status.HTTP_400_BAD_REQUEST, 

1062 detail=ERROR_MESSAGES.DEFAULT('Error deleting file'), 

1063 ) 

1064 else: 

1065 raise HTTPException( 

1066 status_code=status.HTTP_404_NOT_FOUND, 

1067 detail=ERROR_MESSAGES.NOT_FOUND, 

1068 )