Underground-Digital/Workflow-Engine
0
1import logging2import time3 4import click5from celery import shared_task6 7from core.rag.index_processor.index_processor_factory import IndexProcessorFactory8from extensions.ext_database import db9from extensions.ext_storage import storage10from models.dataset import (11 AppDatasetJoin,12 Dataset,13 DatasetProcessRule,14 DatasetQuery,15 Document,16 DocumentSegment,17)18from models.model import UploadFile19 20 21# Add import statement for ValueError22@shared_task(queue="dataset")23def clean_dataset_task(24 dataset_id: str,25 tenant_id: str,26 indexing_technique: str,27 index_struct: str,28 collection_binding_id: str,29 doc_form: str,30):31 """32 Clean dataset when dataset deleted.33 :param dataset_id: dataset id34 :param tenant_id: tenant id35 :param indexing_technique: indexing technique36 :param index_struct: index struct dict37 :param collection_binding_id: collection binding id38 :param doc_form: dataset form39 40 Usage: clean_dataset_task.delay(dataset_id, tenant_id, indexing_technique, index_struct)41 """42 logging.info(click.style("Start clean dataset when dataset deleted: {}".format(dataset_id), fg="green"))43 start_at = time.perf_counter()44 45 try:46 dataset = Dataset(47 id=dataset_id,48 tenant_id=tenant_id,49 indexing_technique=indexing_technique,50 index_struct=index_struct,51 collection_binding_id=collection_binding_id,52 )53 documents = db.session.query(Document).filter(Document.dataset_id == dataset_id).all()54 segments = db.session.query(DocumentSegment).filter(DocumentSegment.dataset_id == dataset_id).all()55 56 if documents is None or len(documents) == 0:57 logging.info(click.style("No documents found for dataset: {}".format(dataset_id), fg="green"))58 else:59 logging.info(click.style("Cleaning documents for dataset: {}".format(dataset_id), fg="green"))60 # Specify the index type before initializing the index processor61 if doc_form is None:62 raise ValueError("Index type must be specified.")63 index_processor = IndexProcessorFactory(doc_form).init_index_processor()64 index_processor.clean(dataset, None)65 66 for document in documents:67 db.session.delete(document)68 69 for segment in segments:70 db.session.delete(segment)71 72 db.session.query(DatasetProcessRule).filter(DatasetProcessRule.dataset_id == dataset_id).delete()73 db.session.query(DatasetQuery).filter(DatasetQuery.dataset_id == dataset_id).delete()74 db.session.query(AppDatasetJoin).filter(AppDatasetJoin.dataset_id == dataset_id).delete()75 76 # delete files77 if documents:78 for document in documents:79 try:80 if document.data_source_type == "upload_file":81 if document.data_source_info:82 data_source_info = document.data_source_info_dict83 if data_source_info and "upload_file_id" in data_source_info:84 file_id = data_source_info["upload_file_id"]85 file = (86 db.session.query(UploadFile)87 .filter(UploadFile.tenant_id == document.tenant_id, UploadFile.id == file_id)88 .first()89 )90 if not file:91 continue92 storage.delete(file.key)93 db.session.delete(file)94 except Exception:95 continue96 97 db.session.commit()98 end_at = time.perf_counter()99 logging.info(100 click.style(101 "Cleaned dataset when dataset deleted: {} latency: {}".format(dataset_id, end_at - start_at), fg="green"102 )103 )104 except Exception:105 logging.exception("Cleaned dataset when dataset deleted failed")106 