from internal.endpoints import exec_endpoint from enrichment import make_enrichment_endpoints, make_post_enrichment_endpoints from airflow.controller import airflow_endpoints from airflow.progress.controller import progress_bar_endpoints from internal.helpers import DataProcessorError import logging from utils.async_task_manager import ThreadWorkerQueue log = logging.getLogger() log.setLevel(logging.INFO) endpoints = {} endpoints.update(make_enrichment_endpoints()) endpoints.update(make_post_enrichment_endpoints()) endpoints.update(airflow_endpoints) endpoints.update(progress_bar_endpoints) # Endpoint for mainly accessing Enrichment process # Called directly by other lambdas/airflow def lambda_handler(event, context): endpoint_name = 'unknown' internal_request = {} try: if 'endpoint' in event: endpoint_name = event["endpoint"] if endpoint_name in endpoints: internal_request = {} result = exec_endpoint(endpoints[endpoint_name], event, internal_request) log.info(f'EP[{endpoint_name}]: {internal_request}') return result else: raise RuntimeError('Unknown endpoint') else: raise RuntimeError('No endpoint defined') except DataProcessorError as e: log.exception(f'EP[{endpoint_name}]: {internal_request}', exc_info=e) return e.args[0] except Exception as e: log.exception(f'EP[{endpoint_name}]: {internal_request}', exc_info=e) return {"error": ', '.join([str(arg) for arg in e.args]), "type": e.__class__.__name__, "info": None} finally: ThreadWorkerQueue.instantiate().wait_for_complete()