Underground-Digital/Workflow-Engine
0
1import json2import logging3import time4 5import click6from celery import shared_task7 8from core.indexing_runner import DocumentIsPausedException9from extensions.ext_database import db10from extensions.ext_storage import storage11from models.dataset import Dataset, ExternalKnowledgeApis12from models.model import UploadFile13from services.external_knowledge_service import ExternalDatasetService14 15 16@shared_task(queue="dataset")17def external_document_indexing_task(18 dataset_id: str, external_knowledge_api_id: str, data_source: dict, process_parameter: dict19):20 """21 Async process document22 :param dataset_id:23 :param external_knowledge_api_id:24 :param data_source:25 :param process_parameter:26 Usage: external_document_indexing_task.delay(dataset_id, document_id)27 """28 start_at = time.perf_counter()29 30 dataset = db.session.query(Dataset).filter(Dataset.id == dataset_id).first()31 if not dataset:32 logging.info(33 click.style("Processed external dataset: {} failed, dataset not exit.".format(dataset_id), fg="red")34 )35 return36 37 # get external api template38 external_knowledge_api = (39 db.session.query(ExternalKnowledgeApis)40 .filter(41 ExternalKnowledgeApis.id == external_knowledge_api_id, ExternalKnowledgeApis.tenant_id == dataset.tenant_id42 )43 .first()44 )45 46 if not external_knowledge_api:47 logging.info(48 click.style(49 "Processed external dataset: {} failed, api template: {} not exit.".format(50 dataset_id, external_knowledge_api_id51 ),52 fg="red",53 )54 )55 return56 files = {}57 if data_source["type"] == "upload_file":58 upload_file_list = data_source["info_list"]["file_info_list"]["file_ids"]59 for file_id in upload_file_list:60 file = (61 db.session.query(UploadFile)62 .filter(UploadFile.tenant_id == dataset.tenant_id, UploadFile.id == file_id)63 .first()64 )65 if file:66 files[file.id] = (file.name, storage.load_once(file.key), file.mime_type)67 try:68 settings = ExternalDatasetService.get_external_knowledge_api_settings(69 json.loads(external_knowledge_api.settings)70 )71 # assemble headers72 headers = ExternalDatasetService.assembling_headers(settings.authorization, settings.headers)73 74 # do http request75 response = ExternalDatasetService.process_external_api(settings, headers, process_parameter, files)76 job_id = response.json().get("job_id")77 if job_id:78 # save job_id to dataset79 dataset.job_id = job_id80 db.session.commit()81 82 end_at = time.perf_counter()83 logging.info(84 click.style(85 "Processed external dataset: {} successful, latency: {}".format(dataset.id, end_at - start_at),86 fg="green",87 )88 )89 except DocumentIsPausedException as ex:90 logging.info(click.style(str(ex), fg="yellow"))91 92 except Exception:93 pass94 