Underground-Digital/Workflow-Engine
0
1import logging2import time3 4import click5from celery import shared_task6from werkzeug.exceptions import NotFound7 8from core.rag.index_processor.index_processor_factory import IndexProcessorFactory9from extensions.ext_database import db10from extensions.ext_redis import redis_client11from models.dataset import Document, DocumentSegment12 13 14@shared_task(queue="dataset")15def remove_document_from_index_task(document_id: str):16 """17 Async Remove document from index18 :param document_id: document id19 20 Usage: remove_document_from_index.delay(document_id)21 """22 logging.info(click.style("Start remove document segments from index: {}".format(document_id), fg="green"))23 start_at = time.perf_counter()24 25 document = db.session.query(Document).filter(Document.id == document_id).first()26 if not document:27 raise NotFound("Document not found")28 29 if document.indexing_status != "completed":30 return31 32 indexing_cache_key = "document_{}_indexing".format(document.id)33 34 try:35 dataset = document.dataset36 37 if not dataset:38 raise Exception("Document has no dataset")39 40 index_processor = IndexProcessorFactory(document.doc_form).init_index_processor()41 42 segments = db.session.query(DocumentSegment).filter(DocumentSegment.document_id == document.id).all()43 index_node_ids = [segment.index_node_id for segment in segments]44 if index_node_ids:45 try:46 index_processor.clean(dataset, index_node_ids)47 except Exception:48 logging.exception(f"clean dataset {dataset.id} from index failed")49 50 end_at = time.perf_counter()51 logging.info(52 click.style(53 "Document removed from index: {} latency: {}".format(document.id, end_at - start_at), fg="green"54 )55 )56 except Exception:57 logging.exception("remove document from index failed")58 if not document.archived:59 document.enabled = True60 db.session.commit()61 finally:62 redis_client.delete(indexing_cache_key)63 