Coverage for open_webui/routers/knowledge.py: 35%
918 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
1from __future__ import annotations
3import asyncio
4import io
5import logging
6import time
7import uuid
8import zipfile
9from typing import List, Optional
10from urllib.parse import quote
12from fastapi import APIRouter, Depends, HTTPException, Query, Request, status
13from fastapi.responses import StreamingResponse
14from open_webui.config import (
15 BYPASS_ADMIN_ACCESS_CONTROL,
16 ENABLE_KNOWLEDGE_FILE_RETENTION,
17 RAG_EMBEDDING_CONTENT_PREFIX,
18)
19from open_webui.constants import ERROR_MESSAGES
20from open_webui.events import EVENTS, publish_event
21from open_webui.internal.db import get_async_session
22from open_webui.models.access_grants import AccessGrants
23from open_webui.models.config import Config
24from open_webui.models.files import FileMetadataResponse, FileModel, FileModelResponse, Files
25from open_webui.models.groups import Groups
26from open_webui.models.knowledge import (
27 KNOWLEDGE_SORTABLE_FIELDS,
28 KnowledgeDirectoryForm,
29 KnowledgeDirectoryModel,
30 KnowledgeFileListResponse,
31 KnowledgeForm,
32 KnowledgeResponse,
33 Knowledges,
34 KnowledgeUserResponse,
35)
36from open_webui.models.models import ModelForm, Models
37from open_webui.retrieval.external import retrieve_external_knowledge, retrieve_external_knowledge_for_connection
38from open_webui.retrieval.vector.async_client import ASYNC_VECTOR_DB_CLIENT
39from open_webui.routers.retrieval import (
40 BatchProcessFilesForm,
41 ProcessFileForm,
42 process_file,
43 process_files_batch,
44)
45from open_webui.storage.provider import Storage
46from open_webui.utils.access_control import filter_allowed_access_grants, has_permission
47from open_webui.utils.access_control.files import has_access_to_file
48from open_webui.utils.auth import get_admin_user, get_verified_user
49from open_webui.utils.json_codec import JSONCodec
50from pydantic import BaseModel
51from sqlalchemy.ext.asyncio import AsyncSession
53log = logging.getLogger(__name__)
55router = APIRouter()
57############################
58# getKnowledgeBases
59############################
61PAGE_ITEM_COUNT = 30
64async def delete_file_resource(file: FileModel, db: AsyncSession) -> bool:
65 try:
66 file_collection = f'file-{file.id}'
67 if await ASYNC_VECTOR_DB_CLIENT.has_collection(collection_name=file_collection):
68 await ASYNC_VECTOR_DB_CLIENT.delete_collection(collection_name=file_collection)
69 except Exception as e:
70 log.debug('This was most likely caused by bypassing embedding processing')
71 log.debug(e)
73 result = await Files.delete_file_by_id(file.id, db=db)
74 if result and file.path:
75 try:
76 await asyncio.to_thread(Storage.delete_file, file.path)
77 except Exception as e:
78 log.debug(e)
80 return result
83############################
84# Knowledge Base Embedding
85############################
87# Knowledge that sits unread serves no one. Let what is
88# stored here find the ones who need it.
89KNOWLEDGE_BASES_COLLECTION = 'knowledge-bases'
92async def embed_knowledge_base_metadata(
93 request: Request,
94 knowledge_base_id: str,
95 name: str,
96 description: str,
97) -> bool:
98 """Generate and store embedding for knowledge base."""
99 try:
100 content = f'{name}\n\n{description}' if description else name
101 embedding = await request.app.state.EMBEDDING_FUNCTION(content, prefix=RAG_EMBEDDING_CONTENT_PREFIX)
102 await ASYNC_VECTOR_DB_CLIENT.upsert(
103 collection_name=KNOWLEDGE_BASES_COLLECTION,
104 items=[
105 {
106 'id': knowledge_base_id,
107 'text': content,
108 'vector': embedding,
109 'metadata': {
110 'knowledge_base_id': knowledge_base_id,
111 },
112 }
113 ],
114 )
115 return True
116 except Exception as e:
117 log.error(f'Failed to embed knowledge base {knowledge_base_id}: {e}')
118 return False
121async def remove_knowledge_base_metadata_embedding(knowledge_base_id: str) -> bool:
122 """Remove knowledge base embedding."""
123 try:
124 await ASYNC_VECTOR_DB_CLIENT.delete(
125 collection_name=KNOWLEDGE_BASES_COLLECTION,
126 ids=[knowledge_base_id],
127 )
128 return True
129 except Exception as e:
130 log.debug('Failed to remove embedding for %s: %s', knowledge_base_id, e)
131 return False
134class KnowledgeAccessResponse(KnowledgeUserResponse):
135 write_access: bool | None = False
138class KnowledgeAccessListResponse(BaseModel):
139 items: list[KnowledgeAccessResponse]
140 total: int
143def is_external_knowledge(knowledge) -> bool:
144 return (knowledge.meta or {}).get('source') == 'external'
147def external_knowledge_error():
148 raise HTTPException(
149 status_code=status.HTTP_400_BAD_REQUEST,
150 detail='External knowledge bases are read-only.',
151 )
154async def _verify_directory_in_knowledge(
155 id: str,
156 directory_id: str | None,
157 db: AsyncSession,
158 detail: str = ERROR_MESSAGES.NOT_FOUND,
159):
160 """Verify a caller-supplied directory belongs to the knowledge base in the URL. Unset means the root level."""
161 if not directory_id:
162 return None
164 directory = await Knowledges.get_directory_by_id(directory_id, db=db)
165 if not directory or directory.knowledge_id != id:
166 raise HTTPException(
167 status_code=status.HTTP_404_NOT_FOUND,
168 detail=detail,
169 )
172@router.get('/', response_model=KnowledgeAccessListResponse)
173async def get_knowledge_bases(
174 page: int | None = 1,
175 user=Depends(get_verified_user),
176 db: AsyncSession = Depends(get_async_session),
177):
178 page = max(page, 1)
179 limit = PAGE_ITEM_COUNT
180 skip = (page - 1) * limit
182 filter = {}
183 groups = await Groups.get_groups_by_member_id(user.id, db=db)
184 user_group_ids = {group.id for group in groups}
186 if not user.role == 'admin' or not BYPASS_ADMIN_ACCESS_CONTROL: 186 ↛ 187line 186 didn't jump to line 187 because the condition on line 186 was never true
187 if groups:
188 filter['group_ids'] = [group.id for group in groups]
190 filter['user_id'] = user.id
192 result = await Knowledges.search_knowledge_bases(user.id, filter=filter, skip=skip, limit=limit, db=db)
194 # Batch-fetch writable knowledge IDs in a single query instead of N has_access calls
195 knowledge_base_ids = [knowledge_base.id for knowledge_base in result.items]
196 writable_knowledge_base_ids = await AccessGrants.get_accessible_resource_ids(
197 user_id=user.id,
198 resource_type='knowledge',
199 resource_ids=knowledge_base_ids,
200 permission='write',
201 user_group_ids=user_group_ids,
202 db=db,
203 )
205 return KnowledgeAccessListResponse(
206 items=[
207 KnowledgeAccessResponse(
208 **knowledge_base.model_dump(),
209 write_access=(
210 user.id == knowledge_base.user_id
211 or (user.role == 'admin' and BYPASS_ADMIN_ACCESS_CONTROL)
212 or knowledge_base.id in writable_knowledge_base_ids
213 ),
214 )
215 for knowledge_base in result.items
216 ],
217 total=result.total,
218 )
221@router.get('/search', response_model=KnowledgeAccessListResponse)
222async def search_knowledge_bases(
223 query: str | None = None,
224 view_option: str | None = None,
225 source: str | None = None,
226 page: int | None = 1,
227 order_by: str | None = None,
228 direction: str | None = None,
229 user=Depends(get_verified_user),
230 db: AsyncSession = Depends(get_async_session),
231):
232 page = max(page, 1)
233 limit = PAGE_ITEM_COUNT
234 skip = (page - 1) * limit
236 filter = {}
237 if query:
238 filter['query'] = query
239 if view_option:
240 filter['view_option'] = view_option
241 if source in {'local', 'external'}: 241 ↛ 242line 241 didn't jump to line 242 because the condition on line 241 was never true
242 filter['source'] = source
243 if order_by in KNOWLEDGE_SORTABLE_FIELDS: 243 ↛ 244line 243 didn't jump to line 244 because the condition on line 243 was never true
244 filter['order_by'] = order_by
245 if direction in {'asc', 'desc'}: 245 ↛ 246line 245 didn't jump to line 246 because the condition on line 245 was never true
246 filter['direction'] = direction
248 groups = await Groups.get_groups_by_member_id(user.id, db=db)
249 user_group_ids = {group.id for group in groups}
251 if not user.role == 'admin' or not BYPASS_ADMIN_ACCESS_CONTROL: 251 ↛ 252line 251 didn't jump to line 252 because the condition on line 251 was never true
252 if groups:
253 filter['group_ids'] = [group.id for group in groups]
255 filter['user_id'] = user.id
257 result = await Knowledges.search_knowledge_bases(user.id, filter=filter, skip=skip, limit=limit, db=db)
259 # Batch-fetch writable knowledge IDs in a single query instead of N has_access calls
260 knowledge_base_ids = [knowledge_base.id for knowledge_base in result.items]
261 writable_knowledge_base_ids = await AccessGrants.get_accessible_resource_ids(
262 user_id=user.id,
263 resource_type='knowledge',
264 resource_ids=knowledge_base_ids,
265 permission='write',
266 user_group_ids=user_group_ids,
267 db=db,
268 )
270 return KnowledgeAccessListResponse(
271 items=[
272 KnowledgeAccessResponse(
273 **knowledge_base.model_dump(),
274 write_access=(
275 user.id == knowledge_base.user_id
276 or (user.role == 'admin' and BYPASS_ADMIN_ACCESS_CONTROL)
277 or knowledge_base.id in writable_knowledge_base_ids
278 ),
279 )
280 for knowledge_base in result.items
281 ],
282 total=result.total,
283 )
286@router.get('/search/files', response_model=KnowledgeFileListResponse)
287async def search_knowledge_files(
288 query: str | None = None,
289 include_content: bool = Query(False, description='Include file content in search (expensive).'),
290 page: int | None = 1,
291 user=Depends(get_verified_user),
292 db: AsyncSession = Depends(get_async_session),
293):
294 page = max(page, 1)
295 limit = PAGE_ITEM_COUNT
296 skip = (page - 1) * limit
298 filter = {}
299 if query:
300 filter['query'] = query
301 if include_content:
302 filter['include_content'] = True
304 groups = await Groups.get_groups_by_member_id(user.id, db=db)
305 if groups: 305 ↛ 306line 305 didn't jump to line 306 because the condition on line 305 was never true
306 filter['group_ids'] = [group.id for group in groups]
308 filter['user_id'] = user.id
310 return await Knowledges.search_knowledge_files(filter=filter, skip=skip, limit=limit, db=db)
313############################
314# CreateNewKnowledge
315############################
318@router.post('/create', response_model=KnowledgeResponse | None)
319async def create_new_knowledge(
320 request: Request,
321 form_data: KnowledgeForm,
322 user=Depends(get_verified_user),
323):
324 # NOTE: We intentionally do NOT use Depends(get_async_session) here.
325 # Database operations (has_permission, filter_allowed_access_grants, insert_new_knowledge) manage their own sessions.
326 # This prevents holding a connection during embed_knowledge_base_metadata()
327 # which makes external embedding API calls (1-5+ seconds).
328 if user.role != 'admin' and not await has_permission( 328 ↛ 331line 328 didn't jump to line 331 because the condition on line 328 was never true
329 user.id, 'workspace.knowledge', await Config.get('user.permissions')
330 ):
331 raise HTTPException(
332 status_code=status.HTTP_401_UNAUTHORIZED,
333 detail=ERROR_MESSAGES.UNAUTHORIZED,
334 )
336 form_data.access_grants = await filter_allowed_access_grants(
337 await Config.get('user.permissions'),
338 user.id,
339 user.role,
340 form_data.access_grants,
341 'sharing.public_knowledge',
342 )
344 knowledge = await Knowledges.insert_new_knowledge(user.id, form_data)
346 if knowledge: 346 ↛ 363line 346 didn't jump to line 363 because the condition on line 346 was always true
347 # Embed knowledge base for semantic search
348 await embed_knowledge_base_metadata(
349 request,
350 knowledge.id,
351 knowledge.name,
352 knowledge.description,
353 )
354 await publish_event(
355 request,
356 EVENTS.KNOWLEDGE_CREATED,
357 actor=user,
358 subject_id=knowledge.id,
359 data={'name': knowledge.name},
360 )
361 return knowledge
362 else:
363 raise HTTPException(
364 status_code=status.HTTP_400_BAD_REQUEST,
365 detail=ERROR_MESSAGES.FILE_EXISTS,
366 )
369############################
370# ReindexKnowledgeFiles
371############################
374@router.post('/reindex', response_model=bool)
375async def reindex_knowledge_files(
376 request: Request,
377 user=Depends(get_verified_user),
378 db: AsyncSession = Depends(get_async_session),
379):
380 if user.role != 'admin': 380 ↛ 381line 380 didn't jump to line 381 because the condition on line 380 was never true
381 raise HTTPException(
382 status_code=status.HTTP_401_UNAUTHORIZED,
383 detail=ERROR_MESSAGES.UNAUTHORIZED,
384 )
386 knowledge_bases = await Knowledges.get_knowledge_bases(db=db)
387 knowledge_base_files = [
388 (knowledge_base, await Knowledges.get_files_by_id(knowledge_base.id, db=db))
389 for knowledge_base in knowledge_bases
390 ]
391 total_files = sum(len(files) for _, files in knowledge_base_files)
392 processed_files = 0
393 failed_files = []
394 start_time = time.monotonic()
396 log.info('Starting reindexing for %s knowledge bases (%s files)', len(knowledge_bases), total_files)
398 for kb_idx, (knowledge_base, files) in enumerate(knowledge_base_files, start=1):
399 try:
400 try:
401 if await ASYNC_VECTOR_DB_CLIENT.has_collection(collection_name=knowledge_base.id): 401 ↛ 402line 401 didn't jump to line 402 because the condition on line 401 was never true
402 await ASYNC_VECTOR_DB_CLIENT.delete_collection(collection_name=knowledge_base.id)
403 except Exception as e:
404 log.error(f'Error deleting collection {knowledge_base.id}: {str(e)}')
405 continue # Skip, don't raise
407 for file in files: 407 ↛ 408line 407 didn't jump to line 408 because the loop on line 407 never started
408 processed_files += 1
409 eta = ''
410 if processed_files > 1:
411 elapsed = time.monotonic() - start_time
412 remaining_files = total_files - processed_files + 1
413 eta = f', ETA: {round(elapsed / (processed_files - 1) * remaining_files)}s'
415 log.info(
416 'Reindexing knowledge base %s/%s file %s/%s%s: %s',
417 kb_idx,
418 len(knowledge_bases),
419 processed_files,
420 total_files,
421 eta,
422 file.filename,
423 )
425 try:
426 # Force the KB add path to use stored SQL content instead of stale file-{id} chunks.
427 # process_file recreates file-{id} only when that stored content exists.
428 file_collection = f'file-{file.id}'
429 if await ASYNC_VECTOR_DB_CLIENT.has_collection(collection_name=file_collection):
430 await ASYNC_VECTOR_DB_CLIENT.delete_collection(collection_name=file_collection)
432 await process_file(
433 request,
434 ProcessFileForm(file_id=file.id, collection_name=knowledge_base.id),
435 user=user,
436 db=db,
437 )
438 except Exception as e:
439 log.error(f'Error processing file {file.filename} (ID: {file.id}): {str(e)}')
440 failed_files.append({'file_id': file.id, 'error': str(e)})
441 continue
443 except Exception as e:
444 log.error(f'Error processing knowledge base {knowledge_base.id}: {str(e)}')
445 # Don't raise, just continue
446 continue
448 if failed_files: 448 ↛ 449line 448 didn't jump to line 449 because the condition on line 448 was never true
449 log.warning(f'Failed to process {len(failed_files)} files')
450 for failed in failed_files:
451 log.warning(f'File ID: {failed["file_id"]}, Error: {failed["error"]}')
453 log.info('Reindexing completed in %ss.', round(time.monotonic() - start_time))
454 await publish_event(
455 request,
456 EVENTS.KNOWLEDGE_REINDEXED,
457 actor=user,
458 subject_id='all',
459 data={'count': len(knowledge_bases)},
460 )
461 return True
464############################
465# ReindexKnowledgeBases
466############################
469@router.post('/metadata/reindex', response_model=dict)
470async def reindex_knowledge_base_metadata_embeddings(
471 request: Request,
472 user=Depends(get_admin_user),
473):
474 """Batch embed all existing knowledge bases. Admin only.
476 NOTE: We intentionally do NOT use Depends(get_async_session) here.
477 This endpoint loops through ALL knowledge bases and calls embed_knowledge_base_metadata()
478 for each one, making N external embedding API calls. Holding a session during
479 this entire operation would exhaust the connection pool.
480 """
481 knowledge_bases = await Knowledges.get_knowledge_bases()
482 log.info('Reindexing embeddings for %s knowledge bases', len(knowledge_bases))
483 try:
484 await ASYNC_VECTOR_DB_CLIENT.delete_collection(collection_name=KNOWLEDGE_BASES_COLLECTION)
485 except Exception as e:
486 log.debug(e)
488 success_count = 0
489 for kb in knowledge_bases:
490 if await embed_knowledge_base_metadata(request, kb.id, kb.name, kb.description):
491 success_count += 1
493 log.info('Embedding reindex complete: %s/%s', success_count, len(knowledge_bases))
494 return {'total': len(knowledge_bases), 'success': success_count}
497############################
498# External Knowledge Sources
499############################
502class ExternalKnowledgeSourceForm(BaseModel):
503 type: str = 'collection'
504 name: str
505 config: Optional[dict] = None
508class ExternalKnowledgeCreateForm(BaseModel):
509 name: str
510 description: str = ''
511 connection_id: str
512 source: ExternalKnowledgeSourceForm
513 access_grants: Optional[list[dict]] = None
516class ExternalKnowledgeSourceCreateForm(BaseModel):
517 name: str
518 description: str = ''
519 connection: ExternalKnowledgeConnectionForm
520 source: ExternalKnowledgeSourceForm
521 access_grants: Optional[list[dict]] = None
522 test_query: str
523 test_count: int = 5
526class ExternalKnowledgeSourceUpdateForm(ExternalKnowledgeSourceCreateForm):
527 pass
530class ExternalKnowledgeSourceTestForm(BaseModel):
531 connection_id: Optional[str] = None
532 connection: ExternalKnowledgeConnectionForm
533 source: ExternalKnowledgeSourceForm
534 query: str
535 count: int = 5
538class ExternalKnowledgeRetrieveTestForm(BaseModel):
539 query: str
540 source: Optional[ExternalKnowledgeSourceForm] = None
541 count: int = 5
544class ExternalKnowledgeConnectionForm(BaseModel):
545 name: str
546 provider: str
547 endpoint: str
548 auth_config: Optional[dict] = None
549 config: Optional[dict] = None
550 capabilities: Optional[dict] = None
551 enabled: bool = True
554class ExternalKnowledgeConnectionListResponse(BaseModel):
555 items: list[dict]
556 total: int
559EXTERNAL_KNOWLEDGE_CONNECTIONS_CONFIG_KEY = 'external_knowledge.connections'
560EXTERNAL_KNOWLEDGE_PROVIDERS = {'qdrant', 'milvus', 'pgvector'}
563def _get_external_connection_provider_and_config(form_data: ExternalKnowledgeConnectionForm) -> tuple[str, dict]:
564 provider = form_data.provider.lower().strip()
565 if provider not in EXTERNAL_KNOWLEDGE_PROVIDERS: 565 ↛ 571line 565 didn't jump to line 571 because the condition on line 565 was always true
566 raise HTTPException(
567 status_code=status.HTTP_400_BAD_REQUEST,
568 detail='Unsupported external knowledge provider.',
569 )
571 if not form_data.name.strip():
572 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Knowledge source name is required.')
574 if not form_data.endpoint.strip():
575 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Knowledge source endpoint is required.')
577 config = form_data.config or {}
578 allowed_config_keys = {'timeout'}
579 if provider == 'milvus':
580 allowed_config_keys.add('db_name')
582 return provider, {key: value for key, value in config.items() if key in allowed_config_keys}
585def _get_external_auth_config(provider: str, incoming: Optional[dict], existing: Optional[dict] = None) -> dict:
586 if provider == 'pgvector':
587 return {}
588 return existing if incoming is None else incoming or {}
591def _get_normalized_external_source(source: ExternalKnowledgeSourceForm, provider: str) -> ExternalKnowledgeSourceForm:
592 source.type = (source.type or 'collection').strip()
593 source.name = source.name.strip()
595 if source.type != 'collection':
596 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Only collection sources are supported.')
597 if not source.name:
598 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Collection name is required.')
600 config = source.config or {}
601 allowed_keys = {'content_field', 'metadata_field', 'document_id_field'}
602 if provider in {'qdrant', 'milvus'}:
603 allowed_keys.add('vector_field')
604 if provider == 'pgvector':
605 allowed_keys.update({'table_name', 'collection_field', 'vector_field'})
607 normalized_config = {
608 key: value.strip() if isinstance(value, str) else value
609 for key, value in config.items()
610 if key in allowed_keys and value is not None and (not isinstance(value, str) or value.strip())
611 }
613 if not normalized_config.get('content_field'):
614 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Content field is required.')
615 if provider in {'milvus', 'pgvector'} and not normalized_config.get('vector_field'):
616 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Vector field is required.')
618 source.config = normalized_config
619 return source
622def _get_sanitized_external_connection(connection: dict) -> dict:
623 sanitized = {**connection}
624 sanitized.pop('auth_config', None)
625 sanitized['auth_configured'] = bool(connection.get('auth_config'))
626 return sanitized
629async def _get_external_connections() -> list[dict]:
630 return await Config.get(EXTERNAL_KNOWLEDGE_CONNECTIONS_CONFIG_KEY, []) or []
633async def _set_external_connections(connections: list[dict]) -> None:
634 await Config.upsert({EXTERNAL_KNOWLEDGE_CONNECTIONS_CONFIG_KEY: connections})
637def _get_external_connection_from_form(
638 form_data: ExternalKnowledgeConnectionForm, user_id: str, id: Optional[str] = None
639) -> dict:
640 provider, config = _get_external_connection_provider_and_config(form_data)
641 now = int(time.time())
642 return {
643 'id': id or str(uuid.uuid4()),
644 'name': form_data.name.strip(),
645 'provider': provider,
646 'endpoint': form_data.endpoint.strip(),
647 'auth_config': _get_external_auth_config(provider, form_data.auth_config),
648 'config': config,
649 'capabilities': form_data.capabilities or {'retrieve': True},
650 'health': None,
651 'enabled': form_data.enabled,
652 'created_by': user_id,
653 'created_at': now,
654 'updated_at': now,
655 }
658def _get_external_connection_update_from_form(
659 form_data: ExternalKnowledgeConnectionForm,
660 existing: dict,
661) -> dict:
662 provider, config = _get_external_connection_provider_and_config(form_data)
663 return {
664 **existing,
665 'name': form_data.name.strip(),
666 'provider': provider,
667 'endpoint': form_data.endpoint.strip(),
668 'auth_config': _get_external_auth_config(provider, form_data.auth_config, existing.get('auth_config')) or {},
669 'config': config,
670 'capabilities': form_data.capabilities or {'retrieve': True},
671 'enabled': form_data.enabled,
672 'updated_at': int(time.time()),
673 }
676async def _get_external_connection_by_id(id: str) -> Optional[dict]:
677 connections = await _get_external_connections()
678 return next((connection for connection in connections if connection.get('id') == id), None)
681async def _get_knowledge_base_count_for_external_connection(
682 connection_id: str, db: Optional[AsyncSession] = None
683) -> int:
684 count = 0
685 for knowledge in await Knowledges.get_knowledge_bases(db=db):
686 if (knowledge.meta or {}).get('external', {}).get('connection_id') == connection_id:
687 count += 1
688 return count
691@router.get('/external/connections', response_model=ExternalKnowledgeConnectionListResponse)
692async def get_external_knowledge_connections(user=Depends(get_admin_user)):
693 connections = [_get_sanitized_external_connection(connection) for connection in await _get_external_connections()]
694 return ExternalKnowledgeConnectionListResponse(items=connections, total=len(connections))
697@router.post('/external/connections', response_model=dict)
698async def create_external_knowledge_connection(
699 request: Request,
700 form_data: ExternalKnowledgeConnectionForm,
701 user=Depends(get_admin_user),
702):
703 connections = await _get_external_connections()
704 connection = _get_external_connection_from_form(form_data, user.id)
705 connections.append(connection)
706 await _set_external_connections(connections)
707 sanitized = _get_sanitized_external_connection(connection)
708 await publish_event(
709 request,
710 EVENTS.KNOWLEDGE_EXTERNAL_CONNECTION_CREATED,
711 actor=user,
712 subject_id=connection.get('id'),
713 data={'name': sanitized.get('name'), 'provider': sanitized.get('provider')},
714 )
715 return sanitized
718@router.get('/external/connections/{id}', response_model=dict)
719async def get_external_knowledge_connection(
720 id: str,
721 user=Depends(get_admin_user),
722):
723 connection = await _get_external_connection_by_id(id)
724 if not connection: 724 ↛ 726line 724 didn't jump to line 726 because the condition on line 724 was always true
725 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
726 return _get_sanitized_external_connection(connection)
729@router.patch('/external/connections/{id}', response_model=dict)
730async def update_external_knowledge_connection(
731 request: Request,
732 id: str,
733 form_data: ExternalKnowledgeConnectionForm,
734 user=Depends(get_admin_user),
735):
736 connections = await _get_external_connections()
737 idx = next((idx for idx, connection in enumerate(connections) if connection.get('id') == id), None)
738 if idx is None: 738 ↛ 741line 738 didn't jump to line 741 because the condition on line 738 was always true
739 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
741 connection = _get_external_connection_update_from_form(form_data, connections[idx])
742 connections[idx] = connection
743 await _set_external_connections(connections)
744 sanitized = _get_sanitized_external_connection(connection)
745 await publish_event(
746 request,
747 EVENTS.KNOWLEDGE_EXTERNAL_CONNECTION_UPDATED,
748 actor=user,
749 subject_id=id,
750 data={'name': sanitized.get('name'), 'provider': sanitized.get('provider')},
751 )
752 return sanitized
755@router.delete('/external/connections/{id}', response_model=bool)
756async def delete_external_knowledge_connection(
757 request: Request,
758 id: str,
759 user=Depends(get_admin_user),
760 db: AsyncSession = Depends(get_async_session),
761):
762 connection = await _get_external_connection_by_id(id)
763 if not connection: 763 ↛ 766line 763 didn't jump to line 766 because the condition on line 763 was always true
764 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
766 if await _get_knowledge_base_count_for_external_connection(id, db=db) > 0:
767 raise HTTPException(
768 status_code=status.HTTP_400_BAD_REQUEST,
769 detail='External connection is still used by knowledge bases.',
770 )
772 connections = [connection for connection in await _get_external_connections() if connection.get('id') != id]
773 await _set_external_connections(connections)
774 await publish_event(
775 request,
776 EVENTS.KNOWLEDGE_EXTERNAL_CONNECTION_DELETED,
777 actor=user,
778 subject_id=id,
779 data={'name': connection.get('name'), 'provider': connection.get('provider')},
780 )
781 return True
784@router.post('/external/connections/{id}/test', response_model=dict)
785async def test_external_knowledge_connection(
786 id: str,
787 user=Depends(get_admin_user),
788):
789 connection = await _get_external_connection_by_id(id)
790 if not connection: 790 ↛ 793line 790 didn't jump to line 793 because the condition on line 790 was always true
791 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
793 health = {
794 'ok': bool(connection.get('enabled') and connection.get('endpoint')),
795 'provider': connection.get('provider'),
796 'checked_at': int(time.time()),
797 }
798 connections = await _get_external_connections()
799 for item in connections:
800 if item.get('id') == id:
801 item['health'] = health
802 item['updated_at'] = int(time.time())
803 break
804 await _set_external_connections(connections)
805 return health
808async def _get_external_source_test_result(
809 request: Request,
810 connection: dict,
811 source: ExternalKnowledgeSourceForm,
812 query: str,
813 count: int,
814 user,
815) -> dict:
816 if not query.strip():
817 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Test query is required.')
819 source = _get_normalized_external_source(source, connection.get('provider'))
820 test_knowledge = KnowledgeResponse(
821 id='external-test',
822 user_id=user.id,
823 name=connection.get('name'),
824 description='',
825 meta={
826 'source': 'external',
827 'read_only': True,
828 'external': {
829 'connection_id': connection.get('id'),
830 'source': source.model_dump(),
831 'provider': connection.get('provider'),
832 'auth_mode': 'service_account',
833 'capabilities': {'retrieve': True},
834 },
835 },
836 access_grants=[],
837 created_at=int(time.time()),
838 updated_at=int(time.time()),
839 )
840 result = await retrieve_external_knowledge_for_connection(
841 request,
842 test_knowledge,
843 connection,
844 [query.strip()],
845 count,
846 user=user,
847 )
848 return {
849 'documents': result.get('documents', [[]])[0],
850 'metadatas': result.get('metadatas', [[]])[0],
851 'distances': result.get('distances', [[]])[0],
852 }
855@router.post('/external/source/test', response_model=dict)
856async def test_external_knowledge_source(
857 request: Request,
858 form_data: ExternalKnowledgeSourceTestForm,
859 user=Depends(get_admin_user),
860):
861 if form_data.connection_id:
862 existing_connection = await _get_external_connection_by_id(form_data.connection_id)
863 if not existing_connection: 863 ↛ 865line 863 didn't jump to line 865 because the condition on line 863 was always true
864 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='External connection not found.')
865 connection = _get_external_connection_update_from_form(form_data.connection, existing_connection)
866 else:
867 connection = _get_external_connection_from_form(form_data.connection, user.id, id='external-test')
869 return await _get_external_source_test_result(
870 request,
871 connection,
872 form_data.source,
873 form_data.query,
874 form_data.count,
875 user,
876 )
879@router.post('/external/connections/{id}/retrieve-test', response_model=dict)
880async def test_external_knowledge_retrieval(
881 request: Request,
882 id: str,
883 form_data: ExternalKnowledgeRetrieveTestForm,
884 user=Depends(get_admin_user),
885):
886 connection = await _get_external_connection_by_id(id)
887 if not connection: 887 ↛ 890line 887 didn't jump to line 890 because the condition on line 887 was always true
888 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
890 source = form_data.source or ExternalKnowledgeSourceForm(name='test', config={'content_field': 'payload.text'})
891 return await _get_external_source_test_result(request, connection, source, form_data.query, form_data.count, user)
894@router.post('/external/knowledge/create', response_model=KnowledgeResponse | None)
895async def create_external_knowledge(
896 request: Request,
897 form_data: ExternalKnowledgeCreateForm,
898 user=Depends(get_admin_user),
899 db: AsyncSession = Depends(get_async_session),
900):
901 connection = await _get_external_connection_by_id(form_data.connection_id)
902 if not connection: 902 ↛ 904line 902 didn't jump to line 904 because the condition on line 902 was always true
903 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
904 if not form_data.name.strip():
905 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Knowledge name is required.')
906 source = _get_normalized_external_source(form_data.source, connection.get('provider'))
908 form_data.access_grants = await filter_allowed_access_grants(
909 await Config.get('user.permissions'),
910 user.id,
911 user.role,
912 form_data.access_grants,
913 'sharing.public_knowledge',
914 )
916 knowledge = await Knowledges.insert_new_knowledge(
917 user.id,
918 KnowledgeForm(
919 name=form_data.name.strip(),
920 description=form_data.description,
921 access_grants=form_data.access_grants,
922 ),
923 db=db,
924 )
925 if not knowledge:
926 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.FILE_EXISTS)
928 meta = {
929 'source': 'external',
930 'read_only': True,
931 'external': {
932 'connection_id': form_data.connection_id,
933 'source': source.model_dump(),
934 'provider': connection.get('provider'),
935 'auth_mode': 'service_account',
936 'capabilities': {'retrieve': True},
937 },
938 }
939 knowledge = await Knowledges.update_knowledge_meta_by_id(knowledge.id, meta, db=db)
940 await embed_knowledge_base_metadata(request, knowledge.id, knowledge.name, knowledge.description)
941 return knowledge
944@router.post('/external/source/create', response_model=KnowledgeResponse | None)
945async def create_external_knowledge_source(
946 request: Request,
947 form_data: ExternalKnowledgeSourceCreateForm,
948 user=Depends(get_admin_user),
949 db: AsyncSession = Depends(get_async_session),
950):
951 if not form_data.name.strip():
952 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Knowledge name is required.')
954 connection = _get_external_connection_from_form(form_data.connection, user.id)
955 source = _get_normalized_external_source(form_data.source, connection.get('provider'))
956 test_result = await _get_external_source_test_result(
957 request,
958 connection,
959 source,
960 form_data.test_query,
961 form_data.test_count,
962 user,
963 )
964 if not test_result.get('documents'):
965 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Test query returned no results.')
967 form_data.access_grants = await filter_allowed_access_grants(
968 await Config.get('user.permissions'),
969 user.id,
970 user.role,
971 form_data.access_grants,
972 'sharing.public_knowledge',
973 )
975 connections = await _get_external_connections()
976 connections.append(connection)
977 await _set_external_connections(connections)
979 knowledge = await Knowledges.insert_new_knowledge(
980 user.id,
981 KnowledgeForm(
982 name=form_data.name.strip(),
983 description=form_data.description,
984 access_grants=form_data.access_grants,
985 ),
986 db=db,
987 )
988 if not knowledge:
989 connections = [item for item in await _get_external_connections() if item.get('id') != connection.get('id')]
990 await _set_external_connections(connections)
991 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.FILE_EXISTS)
993 meta = {
994 'source': 'external',
995 'read_only': True,
996 'external': {
997 'connection_id': connection.get('id'),
998 'source': source.model_dump(),
999 'provider': connection.get('provider'),
1000 'auth_mode': 'service_account',
1001 'capabilities': {'retrieve': True},
1002 },
1003 }
1004 knowledge = await Knowledges.update_knowledge_meta_by_id(knowledge.id, meta, db=db)
1005 await embed_knowledge_base_metadata(request, knowledge.id, knowledge.name, knowledge.description)
1006 return knowledge
1009@router.patch('/external/source/{id}', response_model=KnowledgeResponse | None)
1010async def update_external_knowledge_source(
1011 request: Request,
1012 id: str,
1013 form_data: ExternalKnowledgeSourceUpdateForm,
1014 user=Depends(get_admin_user),
1015 db: AsyncSession = Depends(get_async_session),
1016):
1017 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
1018 if not knowledge or not is_external_knowledge(knowledge): 1018 ↛ 1020line 1018 didn't jump to line 1020 because the condition on line 1018 was always true
1019 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND)
1020 if not form_data.name.strip():
1021 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Knowledge name is required.')
1023 connection_id = (knowledge.meta or {}).get('external', {}).get('connection_id')
1024 connections = await _get_external_connections()
1025 idx = next((idx for idx, connection in enumerate(connections) if connection.get('id') == connection_id), None)
1026 if idx is None:
1027 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail='External connection not found.')
1029 existing_connection = connections[idx]
1030 connection = _get_external_connection_update_from_form(form_data.connection, existing_connection)
1031 source = _get_normalized_external_source(form_data.source, connection.get('provider'))
1032 test_result = await _get_external_source_test_result(
1033 request,
1034 connection,
1035 source,
1036 form_data.test_query,
1037 form_data.test_count,
1038 user,
1039 )
1040 if not test_result.get('documents'):
1041 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail='Test query returned no results.')
1043 form_data.access_grants = await filter_allowed_access_grants(
1044 await Config.get('user.permissions'),
1045 user.id,
1046 user.role,
1047 form_data.access_grants,
1048 'sharing.public_knowledge',
1049 )
1051 connections[idx] = connection
1052 await _set_external_connections(connections)
1054 updated = await Knowledges.update_knowledge_by_id(
1055 id=id,
1056 form_data=KnowledgeForm(
1057 name=form_data.name.strip(),
1058 description=form_data.description,
1059 access_grants=form_data.access_grants,
1060 ),
1061 db=db,
1062 )
1063 if not updated:
1064 connections[idx] = existing_connection
1065 await _set_external_connections(connections)
1066 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1068 meta = {
1069 'source': 'external',
1070 'read_only': True,
1071 'external': {
1072 'connection_id': connection.get('id'),
1073 'source': source.model_dump(),
1074 'provider': connection.get('provider'),
1075 'auth_mode': 'service_account',
1076 'capabilities': {'retrieve': True},
1077 },
1078 }
1079 updated = await Knowledges.update_knowledge_meta_by_id(id, meta, db=db)
1080 if not updated:
1081 connections[idx] = existing_connection
1082 await _set_external_connections(connections)
1083 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=ERROR_MESSAGES.DEFAULT())
1085 await embed_knowledge_base_metadata(request, id, updated.name, updated.description)
1086 return updated
1089############################
1090# GetKnowledgeById
1091############################
1094class KnowledgeFilesResponse(KnowledgeResponse):
1095 files: list[FileMetadataResponse] | None = None
1096 write_access: bool | None = False
1099@router.get('/{id}', response_model=KnowledgeFilesResponse | None)
1100async def get_knowledge_by_id(id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)):
1101 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
1103 if knowledge: 1103 ↛ 1104line 1103 didn't jump to line 1104 because the condition on line 1103 was never true
1104 if (
1105 user.role == 'admin'
1106 or knowledge.user_id == user.id
1107 or await AccessGrants.has_access(
1108 user_id=user.id,
1109 resource_type='knowledge',
1110 resource_id=knowledge.id,
1111 permission='read',
1112 db=db,
1113 )
1114 ):
1115 return KnowledgeFilesResponse(
1116 **knowledge.model_dump(),
1117 write_access=(
1118 user.id == knowledge.user_id
1119 or (user.role == 'admin' and BYPASS_ADMIN_ACCESS_CONTROL)
1120 or await AccessGrants.has_access(
1121 user_id=user.id,
1122 resource_type='knowledge',
1123 resource_id=knowledge.id,
1124 permission='write',
1125 db=db,
1126 )
1127 ),
1128 )
1129 else:
1130 raise HTTPException(
1131 status_code=status.HTTP_401_UNAUTHORIZED,
1132 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
1133 )
1134 else:
1135 raise HTTPException(
1136 status_code=status.HTTP_404_NOT_FOUND,
1137 detail=ERROR_MESSAGES.NOT_FOUND,
1138 )
1141############################
1142# UpdateKnowledgeById
1143############################
1146@router.post('/{id}/update', response_model=KnowledgeFilesResponse | None)
1147async def update_knowledge_by_id(
1148 request: Request,
1149 id: str,
1150 form_data: KnowledgeForm,
1151 user=Depends(get_verified_user),
1152):
1153 # NOTE: We intentionally do NOT use Depends(get_async_session) here.
1154 # Database operations manage their own short-lived sessions internally.
1155 # This prevents holding a connection during embed_knowledge_base_metadata()
1156 # which makes external embedding API calls (1-5+ seconds).
1157 knowledge = await Knowledges.get_knowledge_by_id(id=id)
1158 if not knowledge: 1158 ↛ 1164line 1158 didn't jump to line 1164 because the condition on line 1158 was always true
1159 raise HTTPException(
1160 status_code=status.HTTP_400_BAD_REQUEST,
1161 detail=ERROR_MESSAGES.NOT_FOUND,
1162 )
1163 # Is the user the original creator, in a group with write access, or an admin
1164 if (
1165 knowledge.user_id != user.id
1166 and not await AccessGrants.has_access(
1167 user_id=user.id,
1168 resource_type='knowledge',
1169 resource_id=knowledge.id,
1170 permission='write',
1171 )
1172 and user.role != 'admin'
1173 ):
1174 raise HTTPException(
1175 status_code=status.HTTP_400_BAD_REQUEST,
1176 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
1177 )
1179 form_data.access_grants = await filter_allowed_access_grants(
1180 await Config.get('user.permissions'),
1181 user.id,
1182 user.role,
1183 form_data.access_grants,
1184 'sharing.public_knowledge',
1185 )
1187 knowledge = await Knowledges.update_knowledge_by_id(id=id, form_data=form_data)
1188 if knowledge:
1189 # Re-embed knowledge base for semantic search
1190 await embed_knowledge_base_metadata(
1191 request,
1192 knowledge.id,
1193 knowledge.name,
1194 knowledge.description,
1195 )
1196 response = KnowledgeFilesResponse(
1197 **knowledge.model_dump(),
1198 files=await Knowledges.get_file_metadatas_by_id(knowledge.id),
1199 )
1200 await publish_event(
1201 request,
1202 EVENTS.KNOWLEDGE_UPDATED,
1203 actor=user,
1204 subject_id=knowledge.id,
1205 data={'name': knowledge.name},
1206 )
1207 return response
1208 else:
1209 raise HTTPException(
1210 status_code=status.HTTP_400_BAD_REQUEST,
1211 detail=ERROR_MESSAGES.ID_TAKEN,
1212 )
1215############################
1216# UpdateKnowledgeAccessById
1217############################
1220class KnowledgeAccessGrantsForm(BaseModel):
1221 access_grants: list[dict]
1224@router.post('/{id}/access/update', response_model=KnowledgeFilesResponse | None)
1225async def update_knowledge_access_by_id(
1226 request: Request,
1227 id: str,
1228 form_data: KnowledgeAccessGrantsForm,
1229 user=Depends(get_verified_user),
1230 db: AsyncSession = Depends(get_async_session),
1231):
1232 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
1233 if not knowledge: 1233 ↛ 1239line 1233 didn't jump to line 1239 because the condition on line 1233 was always true
1234 raise HTTPException(
1235 status_code=status.HTTP_404_NOT_FOUND,
1236 detail=ERROR_MESSAGES.NOT_FOUND,
1237 )
1239 if (
1240 knowledge.user_id != user.id
1241 and not await AccessGrants.has_access(
1242 user_id=user.id,
1243 resource_type='knowledge',
1244 resource_id=knowledge.id,
1245 permission='write',
1246 db=db,
1247 )
1248 and user.role != 'admin'
1249 ):
1250 raise HTTPException(
1251 status_code=status.HTTP_400_BAD_REQUEST,
1252 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
1253 )
1255 form_data.access_grants = await filter_allowed_access_grants(
1256 await Config.get('user.permissions'),
1257 user.id,
1258 user.role,
1259 form_data.access_grants,
1260 'sharing.public_knowledge',
1261 )
1263 knowledge.access_grants = await AccessGrants.set_access_grants('knowledge', id, form_data.access_grants, db=db)
1265 response = KnowledgeFilesResponse(
1266 **knowledge.model_dump(),
1267 files=await Knowledges.get_file_metadatas_by_id(id, db=db),
1268 )
1269 await publish_event(
1270 request,
1271 EVENTS.KNOWLEDGE_ACCESS_UPDATED,
1272 actor=user,
1273 subject_id=knowledge.id,
1274 data={'name': knowledge.name},
1275 )
1276 return response
1279############################
1280# GetPendingKnowledgeFiles
1281############################
1284@router.get('/{id}/files/pending')
1285async def get_pending_knowledge_files(
1286 id: str,
1287 stream: bool = Query(False),
1288 user=Depends(get_verified_user),
1289):
1290 """Return files that are being processed for this knowledge base but not yet linked.
1292 After a file is uploaded with ``knowledge_id`` in its metadata, the backend
1293 processes it in a background task before linking it to the ``knowledge_file``
1294 join table. During this window the file is invisible to the normal file
1295 list endpoint. This endpoint exposes those in-flight files so the frontend
1296 can show them with a processing indicator even after a page reload.
1298 When ``stream=true``, returns an SSE stream that polls every 3 seconds
1299 and emits the current pending file list. Closes when no files remain.
1300 """
1301 # NOTE: We intentionally do NOT use Depends(get_async_session) here.
1302 # Database operations manage their own short-lived sessions internally.
1303 # Holding a session here would keep a connection for the entire stream
1304 # (up to an hour) and exhaust the connection pool under concurrent load.
1305 knowledge = await Knowledges.get_knowledge_by_id(id=id)
1306 if not knowledge: 1306 ↛ 1312line 1306 didn't jump to line 1312 because the condition on line 1306 was always true
1307 raise HTTPException(
1308 status_code=status.HTTP_404_NOT_FOUND,
1309 detail=ERROR_MESSAGES.NOT_FOUND,
1310 )
1312 if not (
1313 user.role == 'admin'
1314 or knowledge.user_id == user.id
1315 or await AccessGrants.has_access(
1316 user_id=user.id,
1317 resource_type='knowledge',
1318 resource_id=knowledge.id,
1319 permission='read',
1320 )
1321 ):
1322 raise HTTPException(
1323 status_code=status.HTTP_400_BAD_REQUEST,
1324 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
1325 )
1327 if not stream:
1328 return await Files.get_pending_files_for_knowledge(id)
1330 async def event_stream(knowledge_id: str):
1331 MAX_POLL_DURATION = 3600 # 1 hour max
1332 for _ in range(MAX_POLL_DURATION // 3):
1333 pending = await Files.get_pending_files_for_knowledge(knowledge_id)
1334 data = [f.model_dump() for f in pending]
1335 yield f'data: {JSONCodec.dumps(data)}\n\n'
1336 if len(pending) == 0:
1337 break
1338 await asyncio.sleep(3)
1340 return StreamingResponse(
1341 event_stream(id),
1342 media_type='text/event-stream',
1343 )
1346############################
1347# GetKnowledgeFilesById
1348############################
1351@router.get('/{id}/files', response_model=KnowledgeFileListResponse)
1352async def get_knowledge_files_by_id(
1353 id: str,
1354 query: str | None = None,
1355 include_content: bool = Query(False, description='Include file content in search (expensive).'),
1356 view_option: str | None = None,
1357 order_by: str | None = None,
1358 direction: str | None = None,
1359 directory_id: str | None = Query(None, description='Filter by directory ID. Pass empty string for root.'),
1360 page: int | None = 1,
1361 limit: int | None = Query(None, description='Page size (admin only). Defaults to 30.'),
1362 user=Depends(get_verified_user),
1363 db: AsyncSession = Depends(get_async_session),
1364):
1365 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
1366 if not knowledge: 1366 ↛ 1372line 1366 didn't jump to line 1372 because the condition on line 1366 was always true
1367 raise HTTPException(
1368 status_code=status.HTTP_400_BAD_REQUEST,
1369 detail=ERROR_MESSAGES.NOT_FOUND,
1370 )
1372 if not (
1373 user.role == 'admin'
1374 or knowledge.user_id == user.id
1375 or await AccessGrants.has_access(
1376 user_id=user.id,
1377 resource_type='knowledge',
1378 resource_id=knowledge.id,
1379 permission='read',
1380 db=db,
1381 )
1382 ):
1383 raise HTTPException(
1384 status_code=status.HTTP_400_BAD_REQUEST,
1385 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
1386 )
1388 page = max(page, 1)
1390 # Allow admins to configure page size; non-admins always get the default
1391 if user.role == 'admin' and limit is not None:
1392 limit = max(1, limit)
1393 else:
1394 limit = PAGE_ITEM_COUNT
1395 skip = (page - 1) * limit
1397 filter = {}
1398 if query:
1399 filter['query'] = query
1400 if include_content:
1401 filter['include_content'] = True
1402 if view_option:
1403 filter['view_option'] = view_option
1404 if order_by:
1405 filter['order_by'] = order_by
1406 if direction:
1407 filter['direction'] = direction
1408 # directory_id filtering: present in filter = scope to that directory (None = root)
1409 if directory_id is not None:
1410 filter['directory_id'] = directory_id if directory_id else None
1412 return await Knowledges.search_files_by_id(id, user.id, filter=filter, skip=skip, limit=limit, db=db)
1415############################
1416# AddFileToKnowledge
1417############################
1420class KnowledgeFileIdForm(BaseModel):
1421 file_id: str
1422 directory_id: Optional[str] = None
1425@router.post('/{id}/file/add', response_model=KnowledgeFilesResponse | None)
1426async def add_file_to_knowledge_by_id(
1427 request: Request,
1428 id: str,
1429 form_data: KnowledgeFileIdForm,
1430 user=Depends(get_verified_user),
1431 db: AsyncSession = Depends(get_async_session),
1432):
1433 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
1434 if not knowledge: 1434 ↛ 1439line 1434 didn't jump to line 1439 because the condition on line 1434 was always true
1435 raise HTTPException(
1436 status_code=status.HTTP_400_BAD_REQUEST,
1437 detail=ERROR_MESSAGES.NOT_FOUND,
1438 )
1439 if is_external_knowledge(knowledge):
1440 external_knowledge_error()
1442 if (
1443 knowledge.user_id != user.id
1444 and not await AccessGrants.has_access(
1445 user_id=user.id,
1446 resource_type='knowledge',
1447 resource_id=knowledge.id,
1448 permission='write',
1449 db=db,
1450 )
1451 and user.role != 'admin'
1452 ):
1453 raise HTTPException(
1454 status_code=status.HTTP_400_BAD_REQUEST,
1455 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
1456 )
1458 await _verify_directory_in_knowledge(id, form_data.directory_id, db, detail='Target directory not found.')
1460 file = await Files.get_file_by_id(form_data.file_id, db=db)
1461 if not file:
1462 raise HTTPException(
1463 status_code=status.HTTP_400_BAD_REQUEST,
1464 detail=ERROR_MESSAGES.NOT_FOUND,
1465 )
1466 if not file.data:
1467 raise HTTPException(
1468 status_code=status.HTTP_400_BAD_REQUEST,
1469 detail=ERROR_MESSAGES.FILE_NOT_PROCESSED,
1470 )
1472 # KB write-access alone is not enough — caller must also be able to read the file.
1473 if file.user_id != user.id and user.role != 'admin':
1474 if not await has_access_to_file(file.id, 'read', user, db=db):
1475 raise HTTPException(
1476 status_code=status.HTTP_403_FORBIDDEN,
1477 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
1478 )
1480 # Add content to the vector database
1481 try:
1482 await process_file(
1483 request,
1484 ProcessFileForm(file_id=form_data.file_id, collection_name=id),
1485 user=user,
1486 db=db,
1487 )
1489 # Add file to knowledge base
1490 await Knowledges.add_file_to_knowledge_by_id(
1491 knowledge_id=id,
1492 file_id=form_data.file_id,
1493 user_id=user.id,
1494 directory_id=form_data.directory_id,
1495 db=db,
1496 )
1497 except Exception as e:
1498 log.debug(e)
1499 raise HTTPException(
1500 status_code=status.HTTP_400_BAD_REQUEST,
1501 detail=str(e),
1502 )
1504 if knowledge:
1505 response = KnowledgeFilesResponse(
1506 **knowledge.model_dump(),
1507 files=await Knowledges.get_file_metadatas_by_id(knowledge.id, db=db),
1508 )
1509 await publish_event(
1510 request,
1511 EVENTS.KNOWLEDGE_FILE_ADDED,
1512 actor=user,
1513 subject_id=form_data.file_id,
1514 data={'knowledge_id': knowledge.id, 'directory_id': form_data.directory_id},
1515 )
1516 return response
1517 else:
1518 raise HTTPException(
1519 status_code=status.HTTP_400_BAD_REQUEST,
1520 detail=ERROR_MESSAGES.NOT_FOUND,
1521 )
1524@router.post('/{id}/file/update', response_model=KnowledgeFilesResponse | None)
1525async def update_file_from_knowledge_by_id(
1526 request: Request,
1527 id: str,
1528 form_data: KnowledgeFileIdForm,
1529 user=Depends(get_verified_user),
1530 db: AsyncSession = Depends(get_async_session),
1531):
1532 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
1533 if not knowledge: 1533 ↛ 1538line 1533 didn't jump to line 1538 because the condition on line 1533 was always true
1534 raise HTTPException(
1535 status_code=status.HTTP_400_BAD_REQUEST,
1536 detail=ERROR_MESSAGES.NOT_FOUND,
1537 )
1538 if is_external_knowledge(knowledge):
1539 external_knowledge_error()
1541 if (
1542 knowledge.user_id != user.id
1543 and not await AccessGrants.has_access(
1544 user_id=user.id,
1545 resource_type='knowledge',
1546 resource_id=knowledge.id,
1547 permission='write',
1548 db=db,
1549 )
1550 and user.role != 'admin'
1551 ):
1552 raise HTTPException(
1553 status_code=status.HTTP_400_BAD_REQUEST,
1554 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
1555 )
1557 file = await Files.get_file_by_id(form_data.file_id, db=db)
1558 if not file:
1559 raise HTTPException(
1560 status_code=status.HTTP_400_BAD_REQUEST,
1561 detail=ERROR_MESSAGES.NOT_FOUND,
1562 )
1564 # Validate the file actually belongs to this knowledge base
1565 if not await Knowledges.has_file(knowledge_id=id, file_id=form_data.file_id, db=db):
1566 raise HTTPException(
1567 status_code=status.HTTP_400_BAD_REQUEST,
1568 detail=ERROR_MESSAGES.NOT_FOUND,
1569 )
1571 # Remove content from the vector database
1572 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=knowledge.id, filter={'file_id': form_data.file_id})
1574 # Add content to the vector database
1575 try:
1576 await process_file(
1577 request,
1578 ProcessFileForm(file_id=form_data.file_id, collection_name=id),
1579 user=user,
1580 db=db,
1581 )
1582 except Exception as e:
1583 raise HTTPException(
1584 status_code=status.HTTP_400_BAD_REQUEST,
1585 detail=str(e),
1586 )
1588 if knowledge:
1589 response = KnowledgeFilesResponse(
1590 **knowledge.model_dump(),
1591 files=await Knowledges.get_file_metadatas_by_id(knowledge.id, db=db),
1592 )
1593 await publish_event(
1594 request,
1595 EVENTS.KNOWLEDGE_FILE_UPDATED,
1596 actor=user,
1597 subject_id=form_data.file_id,
1598 data={'knowledge_id': knowledge.id},
1599 )
1600 return response
1601 else:
1602 raise HTTPException(
1603 status_code=status.HTTP_400_BAD_REQUEST,
1604 detail=ERROR_MESSAGES.NOT_FOUND,
1605 )
1608############################
1609# RemoveFileFromKnowledge
1610############################
1613@router.post('/{id}/file/remove', response_model=KnowledgeFilesResponse | None)
1614async def remove_file_from_knowledge_by_id(
1615 request: Request,
1616 id: str,
1617 form_data: KnowledgeFileIdForm,
1618 delete_file: bool = Query(not ENABLE_KNOWLEDGE_FILE_RETENTION),
1619 user=Depends(get_verified_user),
1620 db: AsyncSession = Depends(get_async_session),
1621):
1622 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
1623 if not knowledge: 1623 ↛ 1628line 1623 didn't jump to line 1628 because the condition on line 1623 was always true
1624 raise HTTPException(
1625 status_code=status.HTTP_400_BAD_REQUEST,
1626 detail=ERROR_MESSAGES.NOT_FOUND,
1627 )
1628 if is_external_knowledge(knowledge):
1629 external_knowledge_error()
1631 if (
1632 knowledge.user_id != user.id
1633 and not await AccessGrants.has_access(
1634 user_id=user.id,
1635 resource_type='knowledge',
1636 resource_id=knowledge.id,
1637 permission='write',
1638 db=db,
1639 )
1640 and user.role != 'admin'
1641 ):
1642 raise HTTPException(
1643 status_code=status.HTTP_400_BAD_REQUEST,
1644 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
1645 )
1647 file = await Files.get_file_by_id(form_data.file_id, db=db)
1648 if not file:
1649 raise HTTPException(
1650 status_code=status.HTTP_400_BAD_REQUEST,
1651 detail=ERROR_MESSAGES.NOT_FOUND,
1652 )
1654 # Validate the file actually belongs to this knowledge base
1655 if not await Knowledges.has_file(knowledge_id=id, file_id=form_data.file_id, db=db): 1655 ↛ anywhereline 1655 didn't jump anywhere: it always raised an exception.
1656 raise HTTPException(
1657 status_code=status.HTTP_400_BAD_REQUEST,
1658 detail=ERROR_MESSAGES.NOT_FOUND,
1659 )
1661 await Knowledges.remove_file_from_knowledge_by_id(knowledge_id=id, file_id=form_data.file_id, db=db)
1663 # Remove content from the vector database
1664 try:
1665 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=knowledge.id, filter={'file_id': form_data.file_id})
1666 if file.hash:
1667 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=knowledge.id, filter={'hash': file.hash})
1668 except Exception as e:
1669 log.debug('This was most likely caused by bypassing embedding processing')
1670 log.debug(e)
1671 pass
1673 # Anyone with write permission or higher can delete files
1674 if delete_file and (file.user_id == user.id or user.role == 'admin'):
1675 await delete_file_resource(file, db)
1677 if knowledge:
1678 response = KnowledgeFilesResponse(
1679 **knowledge.model_dump(),
1680 files=await Knowledges.get_file_metadatas_by_id(knowledge.id, db=db),
1681 )
1682 await publish_event(
1683 request,
1684 EVENTS.KNOWLEDGE_FILE_REMOVED,
1685 actor=user,
1686 subject_id=form_data.file_id,
1687 data={'knowledge_id': knowledge.id, 'delete_file': delete_file},
1688 )
1689 return response
1690 else:
1691 raise HTTPException(
1692 status_code=status.HTTP_400_BAD_REQUEST,
1693 detail=ERROR_MESSAGES.NOT_FOUND,
1694 )
1697############################
1698# DeleteKnowledgeById
1699############################
1702@router.delete('/{id}/delete', response_model=bool)
1703async def delete_knowledge_by_id(
1704 request: Request,
1705 id: str,
1706 user=Depends(get_verified_user),
1707 db: AsyncSession = Depends(get_async_session),
1708):
1709 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
1710 if not knowledge: 1710 ↛ 1716line 1710 didn't jump to line 1716 because the condition on line 1710 was always true
1711 raise HTTPException(
1712 status_code=status.HTTP_400_BAD_REQUEST,
1713 detail=ERROR_MESSAGES.NOT_FOUND,
1714 )
1716 if (
1717 knowledge.user_id != user.id
1718 and not await AccessGrants.has_access(
1719 user_id=user.id,
1720 resource_type='knowledge',
1721 resource_id=knowledge.id,
1722 permission='write',
1723 db=db,
1724 )
1725 and user.role != 'admin'
1726 ):
1727 raise HTTPException(
1728 status_code=status.HTTP_400_BAD_REQUEST,
1729 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
1730 )
1732 log.info('Deleting knowledge base: %s (name: %s)', id, knowledge.name)
1734 # Get all models
1735 models = await Models.get_all_models(db=db)
1736 log.info('Found %s models to check for knowledge base %s', len(models), id)
1738 # Update models that reference this knowledge base
1739 for model in models:
1740 if model.meta and hasattr(model.meta, 'knowledge'):
1741 knowledge_list = model.meta.knowledge or []
1742 # Filter out the deleted knowledge base
1743 updated_knowledge = [k for k in knowledge_list if k.get('id') != id]
1745 # If the knowledge list changed, update the model
1746 if len(updated_knowledge) != len(knowledge_list):
1747 log.info('Updating model %s to remove knowledge base %s', model.id, id)
1748 model.meta.knowledge = updated_knowledge
1749 model_form = ModelForm(**model.model_dump())
1750 await Models.update_model_by_id(model.id, model_form, db=db)
1752 # Clean up vector DB
1753 if is_external_knowledge(knowledge):
1754 connection_id = (knowledge.meta or {}).get('external', {}).get('connection_id')
1755 # Connections are admin-owned and shared across knowledge bases
1756 if (
1757 connection_id
1758 and user.role == 'admin'
1759 and await _get_knowledge_base_count_for_external_connection(connection_id, db=db) <= 1
1760 ):
1761 connections = [
1762 connection for connection in await _get_external_connections() if connection.get('id') != connection_id
1763 ]
1764 await _set_external_connections(connections)
1765 else:
1766 try:
1767 await ASYNC_VECTOR_DB_CLIENT.delete_collection(collection_name=id)
1768 except Exception as e:
1769 log.debug(e)
1770 pass
1772 # Remove knowledge base embedding
1773 await remove_knowledge_base_metadata_embedding(id)
1775 result = await Knowledges.delete_knowledge_by_id(id=id, db=db)
1776 if result:
1777 await publish_event(
1778 request,
1779 EVENTS.KNOWLEDGE_DELETED,
1780 actor=user,
1781 subject_id=id,
1782 data={'name': knowledge.name},
1783 )
1784 return result
1787############################
1788# ResetKnowledgeById
1789############################
1792@router.post('/{id}/reset', response_model=KnowledgeResponse | None)
1793async def reset_knowledge_by_id(
1794 request: Request,
1795 id: str,
1796 include_directories: bool = Query(True),
1797 user=Depends(get_verified_user),
1798 db: AsyncSession = Depends(get_async_session),
1799):
1800 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
1801 if not knowledge: 1801 ↛ 1806line 1801 didn't jump to line 1806 because the condition on line 1801 was always true
1802 raise HTTPException(
1803 status_code=status.HTTP_400_BAD_REQUEST,
1804 detail=ERROR_MESSAGES.NOT_FOUND,
1805 )
1806 if is_external_knowledge(knowledge):
1807 external_knowledge_error()
1809 if (
1810 knowledge.user_id != user.id
1811 and not await AccessGrants.has_access(
1812 user_id=user.id,
1813 resource_type='knowledge',
1814 resource_id=knowledge.id,
1815 permission='write',
1816 db=db,
1817 )
1818 and user.role != 'admin'
1819 ):
1820 raise HTTPException(
1821 status_code=status.HTTP_400_BAD_REQUEST,
1822 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
1823 )
1825 files = await Knowledges.get_files_by_id(id, db=db) if not ENABLE_KNOWLEDGE_FILE_RETENTION else []
1827 try:
1828 await ASYNC_VECTOR_DB_CLIENT.delete_collection(collection_name=id)
1829 except Exception as e:
1830 log.debug(e)
1831 pass
1833 for file in files:
1834 if file.user_id == user.id or user.role == 'admin':
1835 await delete_file_resource(file, db)
1837 knowledge = await Knowledges.reset_knowledge_by_id(id=id, include_directories=include_directories, db=db)
1838 if knowledge:
1839 await publish_event(
1840 request,
1841 EVENTS.KNOWLEDGE_RESET,
1842 actor=user,
1843 subject_id=id,
1844 data={'include_directories': include_directories},
1845 )
1846 return knowledge
1849############################
1850# SyncKnowledgeDiff
1851############################
1854class FileManifestEntry(BaseModel):
1855 filename: str # basename: "readme.md"
1856 path: str # relative dir: "docs/api" or "" for root
1857 checksum: str # SHA-256 of raw bytes
1858 size: int
1861class SyncDiffForm(BaseModel):
1862 manifest: list[FileManifestEntry]
1865class SyncDiffResponse(BaseModel):
1866 added: list[dict] # [{filename, path}] — new files
1867 modified: list[dict] # [{filename, path, stale_file_id}] — changed files
1868 deleted: list[dict] # [{file_id, filename}] — files to remove
1869 mkdir: list[str] # directory paths to create
1870 rmdir: list[str] # directory IDs to remove
1871 unmodified_count: int
1872 directory_map: dict[str, str] # existing path → directory ID
1875@router.post('/{id}/sync/diff', response_model=SyncDiffResponse)
1876async def sync_knowledge_diff(
1877 id: str,
1878 form_data: SyncDiffForm,
1879 user=Depends(get_verified_user),
1880 db: AsyncSession = Depends(get_async_session),
1881):
1882 """
1883 Compare a local file manifest against the knowledge base to determine
1884 which files need uploading, removing, and which directories to create/remove.
1885 """
1886 await _verify_knowledge_write_access(id, user, db)
1888 # ── Index existing state ──
1889 knowledge_files = await Knowledges.get_files_with_directory_ids(id, db=db)
1890 existing_directories = await Knowledges.get_all_directories(id, db=db)
1892 # Build directory path lookups
1893 directory_path_by_id: dict[str, str] = {}
1894 directory_id_by_path: dict[str, str] = {}
1895 for directory in existing_directories:
1896 segments = [directory.name]
1897 parent_id = directory.parent_id
1898 while parent_id:
1899 parent = next((d for d in existing_directories if d.id == parent_id), None)
1900 if not parent:
1901 break
1902 segments.insert(0, parent.name)
1903 parent_id = parent.parent_id
1904 full_path = '/'.join(segments)
1905 directory_path_by_id[directory.id] = full_path
1906 directory_id_by_path[full_path] = directory.id
1908 # Index existing files by (path, filename) → {file_id, checksum}
1909 indexed_files: dict[tuple[str, str], dict] = {}
1910 for file_model, directory_id in knowledge_files:
1911 file_path = directory_path_by_id.get(directory_id, '') if directory_id else ''
1912 stored_checksum = (file_model.meta or {}).get('file_hash')
1913 indexed_files[(file_path, file_model.filename)] = {
1914 'file_id': file_model.id,
1915 'checksum': stored_checksum,
1916 }
1918 # ── Diff files ──
1919 added: list[dict] = []
1920 modified: list[dict] = []
1921 deleted: list[dict] = []
1922 unmodified_count = 0
1923 manifest_keys: set[tuple[str, str]] = set()
1925 for entry in form_data.manifest:
1926 key = (entry.path, entry.filename)
1927 manifest_keys.add(key)
1929 if key not in indexed_files:
1930 added.append({'filename': entry.filename, 'path': entry.path})
1931 elif indexed_files[key]['checksum'] != entry.checksum:
1932 modified.append(
1933 {
1934 'filename': entry.filename,
1935 'path': entry.path,
1936 'stale_file_id': indexed_files[key]['file_id'],
1937 }
1938 )
1939 else:
1940 unmodified_count += 1
1942 for key, file_info in indexed_files.items():
1943 if key not in manifest_keys:
1944 deleted.append({'file_id': file_info['file_id'], 'filename': key[1]})
1946 # ── Diff directories ──
1947 required_directory_paths: set[str] = set()
1948 for entry in form_data.manifest:
1949 if entry.path:
1950 segments = entry.path.split('/')
1951 for depth in range(len(segments)):
1952 required_directory_paths.add('/'.join(segments[: depth + 1]))
1954 mkdir = sorted([p for p in required_directory_paths if p not in directory_id_by_path], key=lambda p: p.count('/'))
1956 orphaned_directory_paths = set(directory_id_by_path) - required_directory_paths
1957 rmdir = [directory_id_by_path[p] for p in orphaned_directory_paths]
1959 return SyncDiffResponse(
1960 added=added,
1961 modified=modified,
1962 deleted=deleted,
1963 mkdir=mkdir,
1964 rmdir=rmdir,
1965 unmodified_count=unmodified_count,
1966 directory_map=directory_id_by_path,
1967 )
1970############################
1971# SyncKnowledgeCleanup
1972############################
1975class SyncCleanupForm(BaseModel):
1976 file_ids: list[str] # file IDs to delete
1977 dir_ids: list[str] = [] # directory IDs to rmdir
1980@router.post('/{id}/sync/cleanup')
1981async def sync_knowledge_cleanup(
1982 id: str,
1983 form_data: SyncCleanupForm,
1984 user=Depends(get_verified_user),
1985 db: AsyncSession = Depends(get_async_session),
1986):
1987 """
1988 Remove stale files and orphaned directories from a knowledge base
1989 after an incremental sync.
1990 """
1991 await _verify_knowledge_write_access(id, user, db)
1993 # ── Remove deleted files ──
1994 for file_id in form_data.file_ids:
1995 file = await Files.get_file_by_id(file_id, db=db)
1996 if not file:
1997 continue
1999 # Only clean up files that belong to this knowledge base.
2000 if not await Knowledges.has_file(id, file_id, db=db):
2001 continue
2003 await Knowledges.remove_file_from_knowledge_by_id(id, file_id, db=db)
2005 try:
2006 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=id, filter={'file_id': file_id})
2007 if file.hash:
2008 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=id, filter={'hash': file.hash})
2009 except Exception:
2010 pass
2012 linked_knowledges = await Knowledges.get_knowledges_by_file_id(file_id, db=db)
2013 if (
2014 not ENABLE_KNOWLEDGE_FILE_RETENTION
2015 and not linked_knowledges
2016 and (file.user_id == user.id or user.role == 'admin')
2017 ):
2018 await delete_file_resource(file, db)
2020 # ── Remove orphaned directories (children before parents) ──
2021 for dir_id in reversed(form_data.dir_ids):
2022 # Only delete directories that belong to this knowledge base.
2023 directory = await Knowledges.get_directory_by_id(dir_id, db=db)
2024 if not directory or directory.knowledge_id != id:
2025 continue
2026 await Knowledges.delete_directory(dir_id, move_files_to_parent=False, db=db)
2028 return {'status': True}
2031############################
2032# AddFilesToKnowledge
2033############################
2036@router.post('/{id}/files/batch/add', response_model=KnowledgeFilesResponse | None)
2037async def add_files_to_knowledge_batch(
2038 request: Request,
2039 id: str,
2040 form_data: list[KnowledgeFileIdForm],
2041 user=Depends(get_verified_user),
2042 db: AsyncSession = Depends(get_async_session),
2043):
2044 """
2045 Add multiple files to a knowledge base
2046 """
2047 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
2048 if not knowledge: 2048 ↛ 2053line 2048 didn't jump to line 2053 because the condition on line 2048 was always true
2049 raise HTTPException(
2050 status_code=status.HTTP_400_BAD_REQUEST,
2051 detail=ERROR_MESSAGES.NOT_FOUND,
2052 )
2053 if is_external_knowledge(knowledge):
2054 external_knowledge_error()
2056 if (
2057 knowledge.user_id != user.id
2058 and not await AccessGrants.has_access(
2059 user_id=user.id,
2060 resource_type='knowledge',
2061 resource_id=knowledge.id,
2062 permission='write',
2063 db=db,
2064 )
2065 and user.role != 'admin'
2066 ):
2067 raise HTTPException(
2068 status_code=status.HTTP_400_BAD_REQUEST,
2069 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
2070 )
2072 for directory_id in {form.directory_id for form in form_data if form.directory_id}:
2073 await _verify_directory_in_knowledge(id, directory_id, db, detail='Target directory not found.')
2075 # Batch-fetch all files to avoid N+1 queries
2076 log.info('files/batch/add - %s files', len(form_data))
2077 file_ids = [form.file_id for form in form_data]
2078 files = await Files.get_files_by_ids(file_ids, db=db)
2080 # Verify all requested files were found
2081 found_ids = {file.id for file in files}
2082 missing_ids = [fid for fid in file_ids if fid not in found_ids]
2083 if missing_ids:
2084 raise HTTPException(
2085 status_code=status.HTTP_400_BAD_REQUEST,
2086 detail=f'File {missing_ids[0]} not found',
2087 )
2089 # Per-file read-access check — same gate as the single-file endpoint.
2090 if user.role != 'admin':
2091 for file in files:
2092 if file.user_id != user.id and not await has_access_to_file(file.id, 'read', user, db=db):
2093 raise HTTPException(
2094 status_code=status.HTTP_403_FORBIDDEN,
2095 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
2096 )
2098 # Filter out files already linked to this knowledge base to prevent
2099 # duplicate embeddings in the vector DB (issue #10679).
2100 new_entries = []
2101 for form in form_data:
2102 if not await Knowledges.has_file(knowledge_id=id, file_id=form.file_id, db=db):
2103 new_entries.append(form)
2105 if not new_entries:
2106 return KnowledgeFilesResponse(
2107 **knowledge.model_dump(),
2108 files=await Knowledges.get_file_metadatas_by_id(knowledge.id, db=db),
2109 )
2111 # Narrow the file list to only new files for processing
2112 new_file_ids = {form.file_id for form in new_entries}
2113 files = [f for f in files if f.id in new_file_ids]
2115 # Process files
2116 try:
2117 result = await process_files_batch(
2118 request=request,
2119 form_data=BatchProcessFilesForm(files=files, collection_name=id),
2120 user=user,
2121 db=db,
2122 )
2123 except Exception as e:
2124 log.error(f'add_files_to_knowledge_batch: Exception occurred: {e}', exc_info=True)
2125 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(e))
2127 # Only add files that were successfully processed
2128 successful_file_ids = [r.file_id for r in result.results if r.status == 'completed']
2129 dir_map = {form.file_id: form.directory_id for form in new_entries}
2130 for file_id in successful_file_ids:
2131 await Knowledges.add_file_to_knowledge_by_id(
2132 knowledge_id=id,
2133 file_id=file_id,
2134 user_id=user.id,
2135 directory_id=dir_map.get(file_id),
2136 db=db,
2137 )
2139 # If there were any errors, include them in the response
2140 if result.errors:
2141 error_details = [f'{err.file_id}: {err.error}' for err in result.errors]
2142 return KnowledgeFilesResponse(
2143 **knowledge.model_dump(),
2144 files=await Knowledges.get_file_metadatas_by_id(knowledge.id, db=db),
2145 warnings={
2146 'message': 'Some files failed to process',
2147 'errors': error_details,
2148 },
2149 )
2151 return KnowledgeFilesResponse(
2152 **knowledge.model_dump(),
2153 files=await Knowledges.get_file_metadatas_by_id(knowledge.id, db=db),
2154 )
2157############################
2158# ExportKnowledgeById
2159############################
2162@router.get('/{id}/export')
2163async def export_knowledge_by_id(id: str, user=Depends(get_admin_user), db: AsyncSession = Depends(get_async_session)):
2164 """
2165 Export a knowledge base as a zip file containing .txt files.
2166 Admin only.
2167 """
2169 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
2170 if not knowledge: 2170 ↛ 2175line 2170 didn't jump to line 2175 because the condition on line 2170 was always true
2171 raise HTTPException(
2172 status_code=status.HTTP_404_NOT_FOUND,
2173 detail=ERROR_MESSAGES.NOT_FOUND,
2174 )
2175 if is_external_knowledge(knowledge):
2176 external_knowledge_error()
2178 files = await Knowledges.get_files_by_id(id, db=db)
2180 # Create zip file in memory
2181 zip_buffer = io.BytesIO()
2182 with zipfile.ZipFile(zip_buffer, 'w', zipfile.ZIP_DEFLATED) as zf:
2183 for file in files:
2184 content = file.data.get('content', '') if file.data else ''
2185 if content:
2186 # Use original filename with .txt extension
2187 filename = file.filename
2188 if not filename.endswith('.txt'):
2189 filename = f'{filename}.txt'
2190 zf.writestr(filename, content)
2192 zip_buffer.seek(0)
2194 # Sanitize knowledge name for filename
2195 safe_name = ''.join(c if c.isalnum() or c in ' -_' else '_' for c in knowledge.name)
2196 zip_filename = f'{safe_name}.zip'
2198 return StreamingResponse(
2199 zip_buffer,
2200 media_type='application/zip',
2201 headers={'Content-Disposition': f"attachment; filename*=UTF-8''{quote(zip_filename, safe='')}"},
2202 )
2205############################
2206# Directory endpoints
2207############################
2210class KnowledgeDirectoryCreateForm(BaseModel):
2211 name: str
2212 parent_id: Optional[str] = None
2215class KnowledgeDirectoryUpdateForm(BaseModel):
2216 name: Optional[str] = None
2217 parent_id: Optional[str] = '__unset__'
2220class KnowledgeFileMoveForm(BaseModel):
2221 file_id: str
2222 directory_id: Optional[str] = None
2225async def _verify_knowledge_write_access(id: str, user, db: AsyncSession):
2226 """Verify the user has write access to the knowledge base. Returns the knowledge model."""
2227 knowledge = await Knowledges.get_knowledge_by_id(id=id, db=db)
2228 if not knowledge: 2228 ↛ 2233line 2228 didn't jump to line 2233 because the condition on line 2228 was always true
2229 raise HTTPException(
2230 status_code=status.HTTP_404_NOT_FOUND,
2231 detail=ERROR_MESSAGES.NOT_FOUND,
2232 )
2233 if is_external_knowledge(knowledge):
2234 external_knowledge_error()
2235 if (
2236 knowledge.user_id != user.id
2237 and not await AccessGrants.has_access(
2238 user_id=user.id,
2239 resource_type='knowledge',
2240 resource_id=knowledge.id,
2241 permission='write',
2242 db=db,
2243 )
2244 and user.role != 'admin'
2245 ):
2246 raise HTTPException(
2247 status_code=status.HTTP_403_FORBIDDEN,
2248 detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
2249 )
2250 return knowledge
2253@router.post('/{id}/dirs/create', response_model=KnowledgeDirectoryModel)
2254async def create_knowledge_directory(
2255 request: Request,
2256 id: str,
2257 form_data: KnowledgeDirectoryCreateForm,
2258 user=Depends(get_verified_user),
2259 db: AsyncSession = Depends(get_async_session),
2260):
2261 await _verify_knowledge_write_access(id, user, db)
2263 await _verify_directory_in_knowledge(id, form_data.parent_id, db, detail='Parent directory not found.')
2265 directory = await Knowledges.create_directory(
2266 knowledge_id=id,
2267 name=form_data.name,
2268 user_id=user.id,
2269 parent_id=form_data.parent_id,
2270 db=db,
2271 )
2272 if not directory:
2273 raise HTTPException(
2274 status_code=status.HTTP_400_BAD_REQUEST,
2275 detail='Failed to create directory. A directory with this name may already exist at this level.',
2276 )
2277 await publish_event(
2278 request,
2279 EVENTS.KNOWLEDGE_DIRECTORY_CREATED,
2280 actor=user,
2281 subject_id=directory.id,
2282 data={'knowledge_id': id, 'name': directory.name, 'parent_id': directory.parent_id},
2283 )
2284 return directory
2287@router.post('/{id}/dirs/{dir_id}/update', response_model=KnowledgeDirectoryModel)
2288async def update_knowledge_directory(
2289 request: Request,
2290 id: str,
2291 dir_id: str,
2292 form_data: KnowledgeDirectoryUpdateForm,
2293 user=Depends(get_verified_user),
2294 db: AsyncSession = Depends(get_async_session),
2295):
2296 await _verify_knowledge_write_access(id, user, db)
2297 await _verify_directory_in_knowledge(id, dir_id, db)
2299 # '__unset__' leaves the parent alone, None moves the directory to the root
2300 if form_data.parent_id not in (None, '__unset__'):
2301 await _verify_directory_in_knowledge(id, form_data.parent_id, db, detail='Parent directory not found.')
2303 result = await Knowledges.update_directory(
2304 directory_id=dir_id,
2305 name=form_data.name,
2306 parent_id=form_data.parent_id,
2307 db=db,
2308 )
2309 if not result:
2310 raise HTTPException(
2311 status_code=status.HTTP_400_BAD_REQUEST,
2312 detail='Failed to update directory. This may be caused by a naming conflict or circular move.',
2313 )
2314 await publish_event(
2315 request,
2316 EVENTS.KNOWLEDGE_DIRECTORY_UPDATED,
2317 actor=user,
2318 subject_id=result.id,
2319 data={'knowledge_id': id, 'name': result.name, 'parent_id': result.parent_id},
2320 )
2321 return result
2324@router.delete('/{id}/dirs/{dir_id}/delete')
2325async def delete_knowledge_directory(
2326 request: Request,
2327 id: str,
2328 dir_id: str,
2329 move_files: bool = Query(True, description='If true, move contained files to parent. If false, delete them.'),
2330 user=Depends(get_verified_user),
2331 db: AsyncSession = Depends(get_async_session),
2332):
2333 await _verify_knowledge_write_access(id, user, db)
2334 await _verify_directory_in_knowledge(id, dir_id, db)
2336 # Collect before delete_directory drops the KnowledgeFile rows
2337 files = [] if move_files else await Knowledges.get_files_by_id_and_directory_id(id, dir_id, db=db)
2339 success = await Knowledges.delete_directory(
2340 directory_id=dir_id,
2341 move_files_to_parent=move_files,
2342 db=db,
2343 )
2344 if not success:
2345 raise HTTPException(
2346 status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
2347 detail='Failed to delete directory.',
2348 )
2350 for file in files:
2351 try:
2352 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=id, filter={'file_id': file.id})
2353 if file.hash:
2354 await ASYNC_VECTOR_DB_CLIENT.delete(collection_name=id, filter={'hash': file.hash})
2355 except Exception as e:
2356 log.debug('This was most likely caused by bypassing embedding processing')
2357 log.debug(e)
2359 if (
2360 not ENABLE_KNOWLEDGE_FILE_RETENTION
2361 and not await Knowledges.get_knowledges_by_file_id(file.id, db=db)
2362 and (file.user_id == user.id or user.role == 'admin')
2363 ):
2364 await delete_file_resource(file, db)
2366 await publish_event(
2367 request,
2368 EVENTS.KNOWLEDGE_DIRECTORY_DELETED,
2369 actor=user,
2370 subject_id=dir_id,
2371 data={'knowledge_id': id, 'move_files': move_files},
2372 )
2373 return {'status': True}
2376@router.post('/{id}/file/move')
2377async def move_file_in_knowledge(
2378 request: Request,
2379 id: str,
2380 form_data: KnowledgeFileMoveForm,
2381 user=Depends(get_verified_user),
2382 db: AsyncSession = Depends(get_async_session),
2383):
2384 await _verify_knowledge_write_access(id, user, db)
2386 # Verify file belongs to this knowledge base
2387 if not await Knowledges.has_file(knowledge_id=id, file_id=form_data.file_id, db=db):
2388 raise HTTPException(
2389 status_code=status.HTTP_404_NOT_FOUND,
2390 detail=ERROR_MESSAGES.NOT_FOUND,
2391 )
2393 await _verify_directory_in_knowledge(id, form_data.directory_id, db, detail='Target directory not found.')
2395 success = await Knowledges.move_file_to_directory(
2396 knowledge_id=id,
2397 file_id=form_data.file_id,
2398 directory_id=form_data.directory_id,
2399 db=db,
2400 )
2401 if not success:
2402 raise HTTPException(
2403 status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
2404 detail='Failed to move file.',
2405 )
2406 await publish_event(
2407 request,
2408 EVENTS.KNOWLEDGE_FILE_MOVED,
2409 actor=user,
2410 subject_id=form_data.file_id,
2411 data={'knowledge_id': id, 'directory_id': form_data.directory_id},
2412 )
2413 return {'status': True}