Underground-Digital/Workflow-Engine
0
1import logging2import time3 4import click5from celery import shared_task6from werkzeug.exceptions import NotFound7 8from core.indexing_runner import DocumentIsPausedError, IndexingRunner9from extensions.ext_database import db10from models.dataset import Document11 12 13@shared_task(queue="dataset")14def recover_document_indexing_task(dataset_id: str, document_id: str):15 """16 Async recover document17 :param dataset_id:18 :param document_id:19 20 Usage: recover_document_indexing_task.delay(dataset_id, document_id)21 """22 logging.info(click.style("Recover document: {}".format(document_id), fg="green"))23 start_at = time.perf_counter()24 25 document = db.session.query(Document).filter(Document.id == document_id, Document.dataset_id == dataset_id).first()26 27 if not document:28 raise NotFound("Document not found")29 30 try:31 indexing_runner = IndexingRunner()32 if document.indexing_status in {"waiting", "parsing", "cleaning"}:33 indexing_runner.run([document])34 elif document.indexing_status == "splitting":35 indexing_runner.run_in_splitting_status(document)36 elif document.indexing_status == "indexing":37 indexing_runner.run_in_indexing_status(document)38 end_at = time.perf_counter()39 logging.info(40 click.style("Processed document: {} latency: {}".format(document.id, end_at - start_at), fg="green")41 )42 except DocumentIsPausedError as ex:43 logging.info(click.style(str(ex), fg="yellow"))44 except Exception:45 pass46 