Underground-Digital/Workflow-Engine
0
1import datetime2import logging3import time4 5import click6from celery import shared_task7 8from core.indexing_runner import IndexingRunner9from core.rag.index_processor.index_processor_factory import IndexProcessorFactory10from extensions.ext_database import db11from extensions.ext_redis import redis_client12from models.dataset import Dataset, Document, DocumentSegment13from services.feature_service import FeatureService14 15 16@shared_task(queue="dataset")17def sync_website_document_indexing_task(dataset_id: str, document_id: str):18 """19 Async process document20 :param dataset_id:21 :param document_id:22 23 Usage: sync_website_document_indexing_task.delay(dataset_id, document_id)24 """25 start_at = time.perf_counter()26 27 dataset = db.session.query(Dataset).filter(Dataset.id == dataset_id).first()28 29 sync_indexing_cache_key = "document_{}_is_sync".format(document_id)30 # check document limit31 features = FeatureService.get_features(dataset.tenant_id)32 try:33 if features.billing.enabled:34 vector_space = features.vector_space35 if 0 < vector_space.limit <= vector_space.size:36 raise ValueError(37 "Your total number of documents plus the number of uploads have over the limit of "38 "your subscription."39 )40 except Exception as e:41 document = (42 db.session.query(Document).filter(Document.id == document_id, Document.dataset_id == dataset_id).first()43 )44 if document:45 document.indexing_status = "error"46 document.error = str(e)47 document.stopped_at = datetime.datetime.utcnow()48 db.session.add(document)49 db.session.commit()50 redis_client.delete(sync_indexing_cache_key)51 return52 53 logging.info(click.style("Start sync website document: {}".format(document_id), fg="green"))54 document = db.session.query(Document).filter(Document.id == document_id, Document.dataset_id == dataset_id).first()55 try:56 if document:57 # clean old data58 index_processor = IndexProcessorFactory(document.doc_form).init_index_processor()59 60 segments = db.session.query(DocumentSegment).filter(DocumentSegment.document_id == document_id).all()61 if segments:62 index_node_ids = [segment.index_node_id for segment in segments]63 # delete from vector index64 index_processor.clean(dataset, index_node_ids)65 66 for segment in segments:67 db.session.delete(segment)68 db.session.commit()69 70 document.indexing_status = "parsing"71 document.processing_started_at = datetime.datetime.utcnow()72 db.session.add(document)73 db.session.commit()74 75 indexing_runner = IndexingRunner()76 indexing_runner.run([document])77 redis_client.delete(sync_indexing_cache_key)78 except Exception as ex:79 document.indexing_status = "error"80 document.error = str(ex)81 document.stopped_at = datetime.datetime.utcnow()82 db.session.add(document)83 db.session.commit()84 logging.info(click.style(str(ex), fg="yellow"))85 redis_client.delete(sync_indexing_cache_key)86 pass87 end_at = time.perf_counter()88 logging.info(click.style("Sync document: {} latency: {}".format(document_id, end_at - start_at), fg="green"))89 