Underground-Digital/Workflow-Engine
0
1import concurrent2import logging3from concurrent.futures import ThreadPoolExecutor4from typing import Optional5 6from flask import Flask, current_app7 8from core.app.app_config.entities import ExternalDataVariableEntity9from core.external_data_tool.factory import ExternalDataToolFactory10 11logger = logging.getLogger(__name__)12 13 14class ExternalDataFetch:15 def fetch(16 self,17 tenant_id: str,18 app_id: str,19 external_data_tools: list[ExternalDataVariableEntity],20 inputs: dict,21 query: str,22 ) -> dict:23 """24 Fill in variable inputs from external data tools if exists.25 26 :param tenant_id: workspace id27 :param app_id: app id28 :param external_data_tools: external data tools configs29 :param inputs: the inputs30 :param query: the query31 :return: the filled inputs32 """33 results = {}34 with ThreadPoolExecutor() as executor:35 futures = {}36 for tool in external_data_tools:37 future = executor.submit(38 self._query_external_data_tool,39 current_app._get_current_object(),40 tenant_id,41 app_id,42 tool,43 inputs,44 query,45 )46 47 futures[future] = tool48 49 for future in concurrent.futures.as_completed(futures):50 tool_variable, result = future.result()51 results[tool_variable] = result52 53 inputs.update(results)54 return inputs55 56 def _query_external_data_tool(57 self,58 flask_app: Flask,59 tenant_id: str,60 app_id: str,61 external_data_tool: ExternalDataVariableEntity,62 inputs: dict,63 query: str,64 ) -> tuple[Optional[str], Optional[str]]:65 """66 Query external data tool.67 :param flask_app: flask app68 :param tenant_id: tenant id69 :param app_id: app id70 :param external_data_tool: external data tool71 :param inputs: inputs72 :param query: query73 :return:74 """75 with flask_app.app_context():76 tool_variable = external_data_tool.variable77 tool_type = external_data_tool.type78 tool_config = external_data_tool.config79 80 external_data_tool_factory = ExternalDataToolFactory(81 name=tool_type, tenant_id=tenant_id, app_id=app_id, variable=tool_variable, config=tool_config82 )83 84 # query external data tool85 result = external_data_tool_factory.query(inputs=inputs, query=query)86 87 return tool_variable, result88 