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
« 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
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
51log = logging.getLogger(__name__)
53router = APIRouter()
56from open_webui.utils.access_control.files import has_access_to_file
57from open_webui.utils.json_codec import JSONCodec
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############################
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.
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.
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
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}')
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))
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))
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
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'
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 )
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 )
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')
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 )
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
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
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 )
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)
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 )
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 {}
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
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)
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 {}
346 try:
347 unsanitized_filename = file.filename
348 filename = os.path.basename(unsanitized_filename)
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 ''
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]
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 )
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 )
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 )
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 )
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 )
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)
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 )
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 )
481############################
482# List Files
483############################
486PAGE_SIZE = 50
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
499 result = await Files.get_file_list(user_id=user_id, skip=skip, limit=PAGE_SIZE, db=db)
501 if not content:
502 for file in result.items:
503 if file.data and 'content' in file.data:
504 del file.data['content']
506 return result
509############################
510# Search Files
511############################
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
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 )
542 if not files:
543 raise HTTPException(
544 status_code=status.HTTP_404_NOT_FOUND,
545 detail='No files found matching the pattern.',
546 )
548 if not content:
549 for file in files:
550 if file.data and 'content' in file.data:
551 del file.data['content']
553 return files
556############################
557# Count Files
558############################
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)
570############################
571# Delete All Files
572############################
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 )
600############################
601# Get File By Id
602############################
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)
609 if not file:
610 raise HTTPException(
611 status_code=status.HTTP_404_NOT_FOUND,
612 detail=ERROR_MESSAGES.NOT_FOUND,
613 )
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 )
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)
636 if not file:
637 raise HTTPException(
638 status_code=status.HTTP_404_NOT_FOUND,
639 detail=ERROR_MESSAGES.NOT_FOUND,
640 )
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
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')
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')
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
668 await asyncio.sleep(1)
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 )
683############################
684# Get File Data Content By Id
685############################
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)
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 )
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 )
709############################
710# Update File Data Content By Id
711############################
714class ContentForm(BaseModel):
715 content: str
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)
728 if not file:
729 raise HTTPException(
730 status_code=status.HTTP_404_NOT_FOUND,
731 detail=ERROR_MESSAGES.NOT_FOUND,
732 )
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}')
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 []
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}')
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 )
790############################
791# Get File Content By Id
792############################
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)
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 )
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)
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
821 content_type = file.meta.get('content_type')
822 filename = file.meta.get('name', file.filename)
823 encoded_filename = quote(filename)
824 headers = {}
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}"
835 return FileResponse(file_path, headers=headers, media_type=content_type)
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 )
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)
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 )
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 )
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)
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 )
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)
913 if not file:
914 raise HTTPException(
915 status_code=status.HTTP_404_NOT_FOUND,
916 detail=ERROR_MESSAGES.NOT_FOUND,
917 )
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
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}"}
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)
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
944 # Create a generator that encodes the file content
945 def generator():
946 yield file_content.encode('utf-8')
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 )
960############################
961# Rename File By Id
962############################
965class FileRenameForm(BaseModel):
966 filename: str
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)
979 if not file:
980 raise HTTPException(
981 status_code=status.HTTP_404_NOT_FOUND,
982 detail=ERROR_MESSAGES.NOT_FOUND,
983 )
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 )
1008############################
1009# Delete File By Id
1010############################
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)
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 )
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)
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 )