Underground-Digital/Workflow-Engine
0
1import logging2from argparse import ArgumentTypeError3from datetime import datetime, timezone4 5from flask import request6from flask_login import current_user7from flask_restful import Resource, fields, marshal, marshal_with, reqparse8from sqlalchemy import asc, desc9from transformers.hf_argparser import string_to_bool10from werkzeug.exceptions import Forbidden, NotFound11 12import services13from controllers.console import api14from controllers.console.app.error import (15 ProviderModelCurrentlyNotSupportError,16 ProviderNotInitializeError,17 ProviderQuotaExceededError,18)19from controllers.console.datasets.error import (20 ArchivedDocumentImmutableError,21 DocumentAlreadyFinishedError,22 DocumentIndexingError,23 IndexingEstimateError,24 InvalidActionError,25 InvalidMetadataError,26)27from controllers.console.wraps import (28 account_initialization_required,29 cloud_edition_billing_resource_check,30 setup_required,31)32from core.errors.error import (33 LLMBadRequestError,34 ModelCurrentlyNotSupportError,35 ProviderTokenNotInitError,36 QuotaExceededError,37)38from core.indexing_runner import IndexingRunner39from core.model_manager import ModelManager40from core.model_runtime.entities.model_entities import ModelType41from core.model_runtime.errors.invoke import InvokeAuthorizationError42from core.rag.extractor.entity.extract_setting import ExtractSetting43from extensions.ext_database import db44from extensions.ext_redis import redis_client45from fields.document_fields import (46 dataset_and_document_fields,47 document_fields,48 document_status_fields,49 document_with_segments_fields,50)51from libs.login import login_required52from models import Dataset, DatasetProcessRule, Document, DocumentSegment, UploadFile53from services.dataset_service import DatasetService, DocumentService54from tasks.add_document_to_index_task import add_document_to_index_task55from tasks.remove_document_from_index_task import remove_document_from_index_task56 57 58class DocumentResource(Resource):59 def get_document(self, dataset_id: str, document_id: str) -> Document:60 dataset = DatasetService.get_dataset(dataset_id)61 if not dataset:62 raise NotFound("Dataset not found.")63 64 try:65 DatasetService.check_dataset_permission(dataset, current_user)66 except services.errors.account.NoPermissionError as e:67 raise Forbidden(str(e))68 69 document = DocumentService.get_document(dataset_id, document_id)70 71 if not document:72 raise NotFound("Document not found.")73 74 if document.tenant_id != current_user.current_tenant_id:75 raise Forbidden("No permission.")76 77 return document78 79 def get_batch_documents(self, dataset_id: str, batch: str) -> list[Document]:80 dataset = DatasetService.get_dataset(dataset_id)81 if not dataset:82 raise NotFound("Dataset not found.")83 84 try:85 DatasetService.check_dataset_permission(dataset, current_user)86 except services.errors.account.NoPermissionError as e:87 raise Forbidden(str(e))88 89 documents = DocumentService.get_batch_documents(dataset_id, batch)90 91 if not documents:92 raise NotFound("Documents not found.")93 94 return documents95 96 97class GetProcessRuleApi(Resource):98 @setup_required99 @login_required100 @account_initialization_required101 def get(self):102 req_data = request.args103 104 document_id = req_data.get("document_id")105 106 # get default rules107 mode = DocumentService.DEFAULT_RULES["mode"]108 rules = DocumentService.DEFAULT_RULES["rules"]109 if document_id:110 # get the latest process rule111 document = Document.query.get_or_404(document_id)112 113 dataset = DatasetService.get_dataset(document.dataset_id)114 115 if not dataset:116 raise NotFound("Dataset not found.")117 118 try:119 DatasetService.check_dataset_permission(dataset, current_user)120 except services.errors.account.NoPermissionError as e:121 raise Forbidden(str(e))122 123 # get the latest process rule124 dataset_process_rule = (125 db.session.query(DatasetProcessRule)126 .filter(DatasetProcessRule.dataset_id == document.dataset_id)127 .order_by(DatasetProcessRule.created_at.desc())128 .limit(1)129 .one_or_none()130 )131 if dataset_process_rule:132 mode = dataset_process_rule.mode133 rules = dataset_process_rule.rules_dict134 135 return {"mode": mode, "rules": rules}136 137 138class DatasetDocumentListApi(Resource):139 @setup_required140 @login_required141 @account_initialization_required142 def get(self, dataset_id):143 dataset_id = str(dataset_id)144 page = request.args.get("page", default=1, type=int)145 limit = request.args.get("limit", default=20, type=int)146 search = request.args.get("keyword", default=None, type=str)147 sort = request.args.get("sort", default="-created_at", type=str)148 # "yes", "true", "t", "y", "1" convert to True, while others convert to False.149 try:150 fetch = string_to_bool(request.args.get("fetch", default="false"))151 except (ArgumentTypeError, ValueError, Exception) as e:152 fetch = False153 dataset = DatasetService.get_dataset(dataset_id)154 if not dataset:155 raise NotFound("Dataset not found.")156 157 try:158 DatasetService.check_dataset_permission(dataset, current_user)159 except services.errors.account.NoPermissionError as e:160 raise Forbidden(str(e))161 162 query = Document.query.filter_by(dataset_id=str(dataset_id), tenant_id=current_user.current_tenant_id)163 164 if search:165 search = f"%{search}%"166 query = query.filter(Document.name.like(search))167 168 if sort.startswith("-"):169 sort_logic = desc170 sort = sort[1:]171 else:172 sort_logic = asc173 174 if sort == "hit_count":175 sub_query = (176 db.select(DocumentSegment.document_id, db.func.sum(DocumentSegment.hit_count).label("total_hit_count"))177 .group_by(DocumentSegment.document_id)178 .subquery()179 )180 181 query = query.outerjoin(sub_query, sub_query.c.document_id == Document.id).order_by(182 sort_logic(db.func.coalesce(sub_query.c.total_hit_count, 0)),183 sort_logic(Document.position),184 )185 elif sort == "created_at":186 query = query.order_by(187 sort_logic(Document.created_at),188 sort_logic(Document.position),189 )190 else:191 query = query.order_by(192 desc(Document.created_at),193 desc(Document.position),194 )195 196 paginated_documents = query.paginate(page=page, per_page=limit, max_per_page=100, error_out=False)197 documents = paginated_documents.items198 if fetch:199 for document in documents:200 completed_segments = DocumentSegment.query.filter(201 DocumentSegment.completed_at.isnot(None),202 DocumentSegment.document_id == str(document.id),203 DocumentSegment.status != "re_segment",204 ).count()205 total_segments = DocumentSegment.query.filter(206 DocumentSegment.document_id == str(document.id), DocumentSegment.status != "re_segment"207 ).count()208 document.completed_segments = completed_segments209 document.total_segments = total_segments210 data = marshal(documents, document_with_segments_fields)211 else:212 data = marshal(documents, document_fields)213 response = {214 "data": data,215 "has_more": len(documents) == limit,216 "limit": limit,217 "total": paginated_documents.total,218 "page": page,219 }220 221 return response222 223 documents_and_batch_fields = {"documents": fields.List(fields.Nested(document_fields)), "batch": fields.String}224 225 @setup_required226 @login_required227 @account_initialization_required228 @marshal_with(documents_and_batch_fields)229 @cloud_edition_billing_resource_check("vector_space")230 def post(self, dataset_id):231 dataset_id = str(dataset_id)232 233 dataset = DatasetService.get_dataset(dataset_id)234 235 if not dataset:236 raise NotFound("Dataset not found.")237 238 # The role of the current user in the ta table must be admin, owner, or editor239 if not current_user.is_dataset_editor:240 raise Forbidden()241 242 try:243 DatasetService.check_dataset_permission(dataset, current_user)244 except services.errors.account.NoPermissionError as e:245 raise Forbidden(str(e))246 247 parser = reqparse.RequestParser()248 parser.add_argument(249 "indexing_technique", type=str, choices=Dataset.INDEXING_TECHNIQUE_LIST, nullable=False, location="json"250 )251 parser.add_argument("data_source", type=dict, required=False, location="json")252 parser.add_argument("process_rule", type=dict, required=False, location="json")253 parser.add_argument("duplicate", type=bool, default=True, nullable=False, location="json")254 parser.add_argument("original_document_id", type=str, required=False, location="json")255 parser.add_argument("doc_form", type=str, default="text_model", required=False, nullable=False, location="json")256 parser.add_argument(257 "doc_language", type=str, default="English", required=False, nullable=False, location="json"258 )259 parser.add_argument("retrieval_model", type=dict, required=False, nullable=False, location="json")260 args = parser.parse_args()261 262 if not dataset.indexing_technique and not args["indexing_technique"]:263 raise ValueError("indexing_technique is required.")264 265 # validate args266 DocumentService.document_create_args_validate(args)267 268 try:269 documents, batch = DocumentService.save_document_with_dataset_id(dataset, args, current_user)270 except ProviderTokenNotInitError as ex:271 raise ProviderNotInitializeError(ex.description)272 except QuotaExceededError:273 raise ProviderQuotaExceededError()274 except ModelCurrentlyNotSupportError:275 raise ProviderModelCurrentlyNotSupportError()276 277 return {"documents": documents, "batch": batch}278 279 280class DatasetInitApi(Resource):281 @setup_required282 @login_required283 @account_initialization_required284 @marshal_with(dataset_and_document_fields)285 @cloud_edition_billing_resource_check("vector_space")286 def post(self):287 # The role of the current user in the ta table must be admin, owner, or editor288 if not current_user.is_editor:289 raise Forbidden()290 291 parser = reqparse.RequestParser()292 parser.add_argument(293 "indexing_technique",294 type=str,295 choices=Dataset.INDEXING_TECHNIQUE_LIST,296 required=True,297 nullable=False,298 location="json",299 )300 parser.add_argument("data_source", type=dict, required=True, nullable=True, location="json")301 parser.add_argument("process_rule", type=dict, required=True, nullable=True, location="json")302 parser.add_argument("doc_form", type=str, default="text_model", required=False, nullable=False, location="json")303 parser.add_argument(304 "doc_language", type=str, default="English", required=False, nullable=False, location="json"305 )306 parser.add_argument("retrieval_model", type=dict, required=False, nullable=False, location="json")307 parser.add_argument("embedding_model", type=str, required=False, nullable=True, location="json")308 parser.add_argument("embedding_model_provider", type=str, required=False, nullable=True, location="json")309 args = parser.parse_args()310 311 # The role of the current user in the ta table must be admin, owner, or editor, or dataset_operator312 if not current_user.is_dataset_editor:313 raise Forbidden()314 315 if args["indexing_technique"] == "high_quality":316 if args["embedding_model"] is None or args["embedding_model_provider"] is None:317 raise ValueError("embedding model and embedding model provider are required for high quality indexing.")318 try:319 model_manager = ModelManager()320 model_manager.get_default_model_instance(321 tenant_id=current_user.current_tenant_id, model_type=ModelType.TEXT_EMBEDDING322 )323 except InvokeAuthorizationError:324 raise ProviderNotInitializeError(325 "No Embedding Model available. Please configure a valid provider "326 "in the Settings -> Model Provider."327 )328 except ProviderTokenNotInitError as ex:329 raise ProviderNotInitializeError(ex.description)330 331 # validate args332 DocumentService.document_create_args_validate(args)333 334 try:335 dataset, documents, batch = DocumentService.save_document_without_dataset_id(336 tenant_id=current_user.current_tenant_id, document_data=args, account=current_user337 )338 except ProviderTokenNotInitError as ex:339 raise ProviderNotInitializeError(ex.description)340 except QuotaExceededError:341 raise ProviderQuotaExceededError()342 except ModelCurrentlyNotSupportError:343 raise ProviderModelCurrentlyNotSupportError()344 345 response = {"dataset": dataset, "documents": documents, "batch": batch}346 347 return response348 349 350class DocumentIndexingEstimateApi(DocumentResource):351 @setup_required352 @login_required353 @account_initialization_required354 def get(self, dataset_id, document_id):355 dataset_id = str(dataset_id)356 document_id = str(document_id)357 document = self.get_document(dataset_id, document_id)358 359 if document.indexing_status in {"completed", "error"}:360 raise DocumentAlreadyFinishedError()361 362 data_process_rule = document.dataset_process_rule363 data_process_rule_dict = data_process_rule.to_dict()364 365 response = {"tokens": 0, "total_price": 0, "currency": "USD", "total_segments": 0, "preview": []}366 367 if document.data_source_type == "upload_file":368 data_source_info = document.data_source_info_dict369 if data_source_info and "upload_file_id" in data_source_info:370 file_id = data_source_info["upload_file_id"]371 372 file = (373 db.session.query(UploadFile)374 .filter(UploadFile.tenant_id == document.tenant_id, UploadFile.id == file_id)375 .first()376 )377 378 # raise error if file not found379 if not file:380 raise NotFound("File not found.")381 382 extract_setting = ExtractSetting(383 datasource_type="upload_file", upload_file=file, document_model=document.doc_form384 )385 386 indexing_runner = IndexingRunner()387 388 try:389 response = indexing_runner.indexing_estimate(390 current_user.current_tenant_id,391 [extract_setting],392 data_process_rule_dict,393 document.doc_form,394 "English",395 dataset_id,396 )397 except LLMBadRequestError:398 raise ProviderNotInitializeError(399 "No Embedding Model available. Please configure a valid provider "400 "in the Settings -> Model Provider."401 )402 except ProviderTokenNotInitError as ex:403 raise ProviderNotInitializeError(ex.description)404 except Exception as e:405 raise IndexingEstimateError(str(e))406 407 return response408 409 410class DocumentBatchIndexingEstimateApi(DocumentResource):411 @setup_required412 @login_required413 @account_initialization_required414 def get(self, dataset_id, batch):415 dataset_id = str(dataset_id)416 batch = str(batch)417 documents = self.get_batch_documents(dataset_id, batch)418 response = {"tokens": 0, "total_price": 0, "currency": "USD", "total_segments": 0, "preview": []}419 if not documents:420 return response421 data_process_rule = documents[0].dataset_process_rule422 data_process_rule_dict = data_process_rule.to_dict()423 info_list = []424 extract_settings = []425 for document in documents:426 if document.indexing_status in {"completed", "error"}:427 raise DocumentAlreadyFinishedError()428 data_source_info = document.data_source_info_dict429 # format document files info430 if data_source_info and "upload_file_id" in data_source_info:431 file_id = data_source_info["upload_file_id"]432 info_list.append(file_id)433 # format document notion info434 elif (435 data_source_info and "notion_workspace_id" in data_source_info and "notion_page_id" in data_source_info436 ):437 pages = []438 page = {"page_id": data_source_info["notion_page_id"], "type": data_source_info["type"]}439 pages.append(page)440 notion_info = {"workspace_id": data_source_info["notion_workspace_id"], "pages": pages}441 info_list.append(notion_info)442 443 if document.data_source_type == "upload_file":444 file_id = data_source_info["upload_file_id"]445 file_detail = (446 db.session.query(UploadFile)447 .filter(UploadFile.tenant_id == current_user.current_tenant_id, UploadFile.id == file_id)448 .first()449 )450 451 if file_detail is None:452 raise NotFound("File not found.")453 454 extract_setting = ExtractSetting(455 datasource_type="upload_file", upload_file=file_detail, document_model=document.doc_form456 )457 extract_settings.append(extract_setting)458 459 elif document.data_source_type == "notion_import":460 extract_setting = ExtractSetting(461 datasource_type="notion_import",462 notion_info={463 "notion_workspace_id": data_source_info["notion_workspace_id"],464 "notion_obj_id": data_source_info["notion_page_id"],465 "notion_page_type": data_source_info["type"],466 "tenant_id": current_user.current_tenant_id,467 },468 document_model=document.doc_form,469 )470 extract_settings.append(extract_setting)471 elif document.data_source_type == "website_crawl":472 extract_setting = ExtractSetting(473 datasource_type="website_crawl",474 website_info={475 "provider": data_source_info["provider"],476 "job_id": data_source_info["job_id"],477 "url": data_source_info["url"],478 "tenant_id": current_user.current_tenant_id,479 "mode": data_source_info["mode"],480 "only_main_content": data_source_info["only_main_content"],481 },482 document_model=document.doc_form,483 )484 extract_settings.append(extract_setting)485 486 else:487 raise ValueError("Data source type not support")488 indexing_runner = IndexingRunner()489 try:490 response = indexing_runner.indexing_estimate(491 current_user.current_tenant_id,492 extract_settings,493 data_process_rule_dict,494 document.doc_form,495 "English",496 dataset_id,497 )498 except LLMBadRequestError:499 raise ProviderNotInitializeError(500 "No Embedding Model available. Please configure a valid provider "501 "in the Settings -> Model Provider."502 )503 except ProviderTokenNotInitError as ex:504 raise ProviderNotInitializeError(ex.description)505 except Exception as e:506 raise IndexingEstimateError(str(e))507 return response508 509 510class DocumentBatchIndexingStatusApi(DocumentResource):511 @setup_required512 @login_required513 @account_initialization_required514 def get(self, dataset_id, batch):515 dataset_id = str(dataset_id)516 batch = str(batch)517 documents = self.get_batch_documents(dataset_id, batch)518 documents_status = []519 for document in documents:520 completed_segments = DocumentSegment.query.filter(521 DocumentSegment.completed_at.isnot(None),522 DocumentSegment.document_id == str(document.id),523 DocumentSegment.status != "re_segment",524 ).count()525 total_segments = DocumentSegment.query.filter(526 DocumentSegment.document_id == str(document.id), DocumentSegment.status != "re_segment"527 ).count()528 document.completed_segments = completed_segments529 document.total_segments = total_segments530 if document.is_paused:531 document.indexing_status = "paused"532 documents_status.append(marshal(document, document_status_fields))533 data = {"data": documents_status}534 return data535 536 537class DocumentIndexingStatusApi(DocumentResource):538 @setup_required539 @login_required540 @account_initialization_required541 def get(self, dataset_id, document_id):542 dataset_id = str(dataset_id)543 document_id = str(document_id)544 document = self.get_document(dataset_id, document_id)545 546 completed_segments = DocumentSegment.query.filter(547 DocumentSegment.completed_at.isnot(None),548 DocumentSegment.document_id == str(document_id),549 DocumentSegment.status != "re_segment",550 ).count()551 total_segments = DocumentSegment.query.filter(552 DocumentSegment.document_id == str(document_id), DocumentSegment.status != "re_segment"553 ).count()554 555 document.completed_segments = completed_segments556 document.total_segments = total_segments557 if document.is_paused:558 document.indexing_status = "paused"559 return marshal(document, document_status_fields)560 561 562class DocumentDetailApi(DocumentResource):563 METADATA_CHOICES = {"all", "only", "without"}564 565 @setup_required566 @login_required567 @account_initialization_required568 def get(self, dataset_id, document_id):569 dataset_id = str(dataset_id)570 document_id = str(document_id)571 document = self.get_document(dataset_id, document_id)572 573 metadata = request.args.get("metadata", "all")574 if metadata not in self.METADATA_CHOICES:575 raise InvalidMetadataError(f"Invalid metadata value: {metadata}")576 577 if metadata == "only":578 response = {"id": document.id, "doc_type": document.doc_type, "doc_metadata": document.doc_metadata}579 elif metadata == "without":580 process_rules = DatasetService.get_process_rules(dataset_id)581 data_source_info = document.data_source_detail_dict582 response = {583 "id": document.id,584 "position": document.position,585 "data_source_type": document.data_source_type,586 "data_source_info": data_source_info,587 "dataset_process_rule_id": document.dataset_process_rule_id,588 "dataset_process_rule": process_rules,589 "name": document.name,590 "created_from": document.created_from,591 "created_by": document.created_by,592 "created_at": document.created_at.timestamp(),593 "tokens": document.tokens,594 "indexing_status": document.indexing_status,595 "completed_at": int(document.completed_at.timestamp()) if document.completed_at else None,596 "updated_at": int(document.updated_at.timestamp()) if document.updated_at else None,597 "indexing_latency": document.indexing_latency,598 "error": document.error,599 "enabled": document.enabled,600 "disabled_at": int(document.disabled_at.timestamp()) if document.disabled_at else None,601 "disabled_by": document.disabled_by,602 "archived": document.archived,603 "segment_count": document.segment_count,604 "average_segment_length": document.average_segment_length,605 "hit_count": document.hit_count,606 "display_status": document.display_status,607 "doc_form": document.doc_form,608 "doc_language": document.doc_language,609 }610 else:611 process_rules = DatasetService.get_process_rules(dataset_id)612 data_source_info = document.data_source_detail_dict613 response = {614 "id": document.id,615 "position": document.position,616 "data_source_type": document.data_source_type,617 "data_source_info": data_source_info,618 "dataset_process_rule_id": document.dataset_process_rule_id,619 "dataset_process_rule": process_rules,620 "name": document.name,621 "created_from": document.created_from,622 "created_by": document.created_by,623 "created_at": document.created_at.timestamp(),624 "tokens": document.tokens,625 "indexing_status": document.indexing_status,626 "completed_at": int(document.completed_at.timestamp()) if document.completed_at else None,627 "updated_at": int(document.updated_at.timestamp()) if document.updated_at else None,628 "indexing_latency": document.indexing_latency,629 "error": document.error,630 "enabled": document.enabled,631 "disabled_at": int(document.disabled_at.timestamp()) if document.disabled_at else None,632 "disabled_by": document.disabled_by,633 "archived": document.archived,634 "doc_type": document.doc_type,635 "doc_metadata": document.doc_metadata,636 "segment_count": document.segment_count,637 "average_segment_length": document.average_segment_length,638 "hit_count": document.hit_count,639 "display_status": document.display_status,640 "doc_form": document.doc_form,641 "doc_language": document.doc_language,642 }643 644 return response, 200645 646 647class DocumentProcessingApi(DocumentResource):648 @setup_required649 @login_required650 @account_initialization_required651 def patch(self, dataset_id, document_id, action):652 dataset_id = str(dataset_id)653 document_id = str(document_id)654 document = self.get_document(dataset_id, document_id)655 656 # The role of the current user in the ta table must be admin, owner, or editor657 if not current_user.is_editor:658 raise Forbidden()659 660 if action == "pause":661 if document.indexing_status != "indexing":662 raise InvalidActionError("Document not in indexing state.")663 664 document.paused_by = current_user.id665 document.paused_at = datetime.now(timezone.utc).replace(tzinfo=None)666 document.is_paused = True667 db.session.commit()668 669 elif action == "resume":670 if document.indexing_status not in {"paused", "error"}:671 raise InvalidActionError("Document not in paused or error state.")672 673 document.paused_by = None674 document.paused_at = None675 document.is_paused = False676 db.session.commit()677 else:678 raise InvalidActionError()679 680 return {"result": "success"}, 200681 682 683class DocumentDeleteApi(DocumentResource):684 @setup_required685 @login_required686 @account_initialization_required687 def delete(self, dataset_id, document_id):688 dataset_id = str(dataset_id)689 document_id = str(document_id)690 dataset = DatasetService.get_dataset(dataset_id)691 if dataset is None:692 raise NotFound("Dataset not found.")693 # check user's model setting694 DatasetService.check_dataset_model_setting(dataset)695 696 document = self.get_document(dataset_id, document_id)697 698 try:699 DocumentService.delete_document(document)700 except services.errors.document.DocumentIndexingError:701 raise DocumentIndexingError("Cannot delete document during indexing.")702 703 return {"result": "success"}, 204704 705 706class DocumentMetadataApi(DocumentResource):707 @setup_required708 @login_required709 @account_initialization_required710 def put(self, dataset_id, document_id):711 dataset_id = str(dataset_id)712 document_id = str(document_id)713 document = self.get_document(dataset_id, document_id)714 715 req_data = request.get_json()716 717 doc_type = req_data.get("doc_type")718 doc_metadata = req_data.get("doc_metadata")719 720 # The role of the current user in the ta table must be admin, owner, or editor721 if not current_user.is_editor:722 raise Forbidden()723 724 if doc_type is None or doc_metadata is None:725 raise ValueError("Both doc_type and doc_metadata must be provided.")726 727 if doc_type not in DocumentService.DOCUMENT_METADATA_SCHEMA:728 raise ValueError("Invalid doc_type.")729 730 if not isinstance(doc_metadata, dict):731 raise ValueError("doc_metadata must be a dictionary.")732 733 metadata_schema = DocumentService.DOCUMENT_METADATA_SCHEMA[doc_type]734 735 document.doc_metadata = {}736 if doc_type == "others":737 document.doc_metadata = doc_metadata738 else:739 for key, value_type in metadata_schema.items():740 value = doc_metadata.get(key)741 if value is not None and isinstance(value, value_type):742 document.doc_metadata[key] = value743 744 document.doc_type = doc_type745 document.updated_at = datetime.now(timezone.utc).replace(tzinfo=None)746 db.session.commit()747 748 return {"result": "success", "message": "Document metadata updated."}, 200749 750 751class DocumentStatusApi(DocumentResource):752 @setup_required753 @login_required754 @account_initialization_required755 @cloud_edition_billing_resource_check("vector_space")756 def patch(self, dataset_id, document_id, action):757 dataset_id = str(dataset_id)758 document_id = str(document_id)759 dataset = DatasetService.get_dataset(dataset_id)760 if dataset is None:761 raise NotFound("Dataset not found.")762 763 # The role of the current user in the ta table must be admin, owner, or editor764 if not current_user.is_dataset_editor:765 raise Forbidden()766 767 # check user's model setting768 DatasetService.check_dataset_model_setting(dataset)769 770 # check user's permission771 DatasetService.check_dataset_permission(dataset, current_user)772 773 document = self.get_document(dataset_id, document_id)774 775 indexing_cache_key = "document_{}_indexing".format(document.id)776 cache_result = redis_client.get(indexing_cache_key)777 if cache_result is not None:778 raise InvalidActionError("Document is being indexed, please try again later")779 780 if action == "enable":781 if document.enabled:782 raise InvalidActionError("Document already enabled.")783 784 document.enabled = True785 document.disabled_at = None786 document.disabled_by = None787 document.updated_at = datetime.now(timezone.utc).replace(tzinfo=None)788 db.session.commit()789 790 # Set cache to prevent indexing the same document multiple times791 redis_client.setex(indexing_cache_key, 600, 1)792 793 add_document_to_index_task.delay(document_id)794 795 return {"result": "success"}, 200796 797 elif action == "disable":798 if not document.completed_at or document.indexing_status != "completed":799 raise InvalidActionError("Document is not completed.")800 if not document.enabled:801 raise InvalidActionError("Document already disabled.")802 803 document.enabled = False804 document.disabled_at = datetime.now(timezone.utc).replace(tzinfo=None)805 document.disabled_by = current_user.id806 document.updated_at = datetime.now(timezone.utc).replace(tzinfo=None)807 db.session.commit()808 809 # Set cache to prevent indexing the same document multiple times810 redis_client.setex(indexing_cache_key, 600, 1)811 812 remove_document_from_index_task.delay(document_id)813 814 return {"result": "success"}, 200815 816 elif action == "archive":817 if document.archived:818 raise InvalidActionError("Document already archived.")819 820 document.archived = True821 document.archived_at = datetime.now(timezone.utc).replace(tzinfo=None)822 document.archived_by = current_user.id823 document.updated_at = datetime.now(timezone.utc).replace(tzinfo=None)824 db.session.commit()825 826 if document.enabled:827 # Set cache to prevent indexing the same document multiple times828 redis_client.setex(indexing_cache_key, 600, 1)829 830 remove_document_from_index_task.delay(document_id)831 832 return {"result": "success"}, 200833 elif action == "un_archive":834 if not document.archived:835 raise InvalidActionError("Document is not archived.")836 837 document.archived = False838 document.archived_at = None839 document.archived_by = None840 document.updated_at = datetime.now(timezone.utc).replace(tzinfo=None)841 db.session.commit()842 843 # Set cache to prevent indexing the same document multiple times844 redis_client.setex(indexing_cache_key, 600, 1)845 846 add_document_to_index_task.delay(document_id)847 848 return {"result": "success"}, 200849 else:850 raise InvalidActionError()851 852 853class DocumentPauseApi(DocumentResource):854 @setup_required855 @login_required856 @account_initialization_required857 def patch(self, dataset_id, document_id):858 """pause document."""859 dataset_id = str(dataset_id)860 document_id = str(document_id)861 862 dataset = DatasetService.get_dataset(dataset_id)863 if not dataset:864 raise NotFound("Dataset not found.")865 866 document = DocumentService.get_document(dataset.id, document_id)867 868 # 404 if document not found869 if document is None:870 raise NotFound("Document Not Exists.")871 872 # 403 if document is archived873 if DocumentService.check_archived(document):874 raise ArchivedDocumentImmutableError()875 876 try:877 # pause document878 DocumentService.pause_document(document)879 except services.errors.document.DocumentIndexingError:880 raise DocumentIndexingError("Cannot pause completed document.")881 882 return {"result": "success"}, 204883 884 885class DocumentRecoverApi(DocumentResource):886 @setup_required887 @login_required888 @account_initialization_required889 def patch(self, dataset_id, document_id):890 """recover document."""891 dataset_id = str(dataset_id)892 document_id = str(document_id)893 dataset = DatasetService.get_dataset(dataset_id)894 if not dataset:895 raise NotFound("Dataset not found.")896 document = DocumentService.get_document(dataset.id, document_id)897 898 # 404 if document not found899 if document is None:900 raise NotFound("Document Not Exists.")901 902 # 403 if document is archived903 if DocumentService.check_archived(document):904 raise ArchivedDocumentImmutableError()905 try:906 # pause document907 DocumentService.recover_document(document)908 except services.errors.document.DocumentIndexingError:909 raise DocumentIndexingError("Document is not in paused status.")910 911 return {"result": "success"}, 204912 913 914class DocumentRetryApi(DocumentResource):915 @setup_required916 @login_required917 @account_initialization_required918 def post(self, dataset_id):919 """retry document."""920 921 parser = reqparse.RequestParser()922 parser.add_argument("document_ids", type=list, required=True, nullable=False, location="json")923 args = parser.parse_args()924 dataset_id = str(dataset_id)925 dataset = DatasetService.get_dataset(dataset_id)926 retry_documents = []927 if not dataset:928 raise NotFound("Dataset not found.")929 for document_id in args["document_ids"]:930 try:931 document_id = str(document_id)932 933 document = DocumentService.get_document(dataset.id, document_id)934 935 # 404 if document not found936 if document is None:937 raise NotFound("Document Not Exists.")938 939 # 403 if document is archived940 if DocumentService.check_archived(document):941 raise ArchivedDocumentImmutableError()942 943 # 400 if document is completed944 if document.indexing_status == "completed":945 raise DocumentAlreadyFinishedError()946 retry_documents.append(document)947 except Exception as e:948 logging.error(f"Document {document_id} retry failed: {str(e)}")949 continue950 # retry document951 DocumentService.retry_document(dataset_id, retry_documents)952 953 return {"result": "success"}, 204954 955 956class DocumentRenameApi(DocumentResource):957 @setup_required958 @login_required959 @account_initialization_required960 @marshal_with(document_fields)961 def post(self, dataset_id, document_id):962 # The role of the current user in the ta table must be admin, owner, editor, or dataset_operator963 if not current_user.is_dataset_editor:964 raise Forbidden()965 dataset = DatasetService.get_dataset(dataset_id)966 DatasetService.check_dataset_operator_permission(current_user, dataset)967 parser = reqparse.RequestParser()968 parser.add_argument("name", type=str, required=True, nullable=False, location="json")969 args = parser.parse_args()970 971 try:972 document = DocumentService.rename_document(dataset_id, document_id, args["name"])973 except services.errors.document.DocumentIndexingError:974 raise DocumentIndexingError("Cannot delete document during indexing.")975 976 return document977 978 979class WebsiteDocumentSyncApi(DocumentResource):980 @setup_required981 @login_required982 @account_initialization_required983 def get(self, dataset_id, document_id):984 """sync website document."""985 dataset_id = str(dataset_id)986 dataset = DatasetService.get_dataset(dataset_id)987 if not dataset:988 raise NotFound("Dataset not found.")989 document_id = str(document_id)990 document = DocumentService.get_document(dataset.id, document_id)991 if not document:992 raise NotFound("Document not found.")993 if document.tenant_id != current_user.current_tenant_id:994 raise Forbidden("No permission.")995 if document.data_source_type != "website_crawl":996 raise ValueError("Document is not a website document.")997 # 403 if document is archived998 if DocumentService.check_archived(document):999 raise ArchivedDocumentImmutableError()1000 # sync document1001 DocumentService.sync_website_document(dataset_id, document)1002 1003 return {"result": "success"}, 2001004 1005 1006api.add_resource(GetProcessRuleApi, "/datasets/process-rule")1007api.add_resource(DatasetDocumentListApi, "/datasets/<uuid:dataset_id>/documents")1008api.add_resource(DatasetInitApi, "/datasets/init")1009api.add_resource(1010 DocumentIndexingEstimateApi, "/datasets/<uuid:dataset_id>/documents/<uuid:document_id>/indexing-estimate"1011)1012api.add_resource(DocumentBatchIndexingEstimateApi, "/datasets/<uuid:dataset_id>/batch/<string:batch>/indexing-estimate")1013api.add_resource(DocumentBatchIndexingStatusApi, "/datasets/<uuid:dataset_id>/batch/<string:batch>/indexing-status")1014api.add_resource(DocumentIndexingStatusApi, "/datasets/<uuid:dataset_id>/documents/<uuid:document_id>/indexing-status")1015api.add_resource(DocumentDetailApi, "/datasets/<uuid:dataset_id>/documents/<uuid:document_id>")1016api.add_resource(1017 DocumentProcessingApi, "/datasets/<uuid:dataset_id>/documents/<uuid:document_id>/processing/<string:action>"1018)1019api.add_resource(DocumentDeleteApi, "/datasets/<uuid:dataset_id>/documents/<uuid:document_id>")1020api.add_resource(DocumentMetadataApi, "/datasets/<uuid:dataset_id>/documents/<uuid:document_id>/metadata")1021api.add_resource(DocumentStatusApi, "/datasets/<uuid:dataset_id>/documents/<uuid:document_id>/status/<string:action>")1022api.add_resource(DocumentPauseApi, "/datasets/<uuid:dataset_id>/documents/<uuid:document_id>/processing/pause")1023api.add_resource(DocumentRecoverApi, "/datasets/<uuid:dataset_id>/documents/<uuid:document_id>/processing/resume")1024api.add_resource(DocumentRetryApi, "/datasets/<uuid:dataset_id>/retry")1025api.add_resource(DocumentRenameApi, "/datasets/<uuid:dataset_id>/documents/<uuid:document_id>/rename")1026 1027api.add_resource(WebsiteDocumentSyncApi, "/datasets/<uuid:dataset_id>/documents/<uuid:document_id>/website-sync")1028 