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_redis import redis_client10from models.dataset import Dataset, Document11 12 13@shared_task(queue="dataset")14def delete_segment_from_index_task(segment_id: str, index_node_id: str, dataset_id: str, document_id: str):15 """16 Async Remove segment from index17 :param segment_id:18 :param index_node_id:19 :param dataset_id:20 :param document_id:21 22 Usage: delete_segment_from_index_task.delay(segment_id)23 """24 logging.info(click.style("Start delete segment from index: {}".format(segment_id), fg="green"))25 start_at = time.perf_counter()26 indexing_cache_key = "segment_{}_delete_indexing".format(segment_id)27 try:28 dataset = db.session.query(Dataset).filter(Dataset.id == dataset_id).first()29 if not dataset:30 logging.info(click.style("Segment {} has no dataset, pass.".format(segment_id), fg="cyan"))31 return32 33 dataset_document = db.session.query(Document).filter(Document.id == document_id).first()34 if not dataset_document:35 logging.info(click.style("Segment {} has no document, pass.".format(segment_id), fg="cyan"))36 return37 38 if not dataset_document.enabled or dataset_document.archived or dataset_document.indexing_status != "completed":39 logging.info(click.style("Segment {} document status is invalid, pass.".format(segment_id), fg="cyan"))40 return41 42 index_type = dataset_document.doc_form43 index_processor = IndexProcessorFactory(index_type).init_index_processor()44 index_processor.clean(dataset, [index_node_id])45 46 end_at = time.perf_counter()47 logging.info(48 click.style("Segment deleted from index: {} latency: {}".format(segment_id, end_at - start_at), fg="green")49 )50 except Exception:51 logging.exception("delete segment from index failed")52 finally:53 redis_client.delete(indexing_cache_key)54 