import logging from airflow.progress.controller import get_enrichment_progress from internal.commons import DP_URL from internal.queries import http_query from utils import appsync_communication as appsync from utils.data_model_utils import change_collection_status from utils.module_tools import class_dict from .enrich_age import EnrichAge from .enrich_feature_eng import EnrichFeatureEngineering from .enrich_geo_pelias import EnrichGeographicsPelias from .enrich_ml import EnrichMachineLearningClusters from .enrich_name import EnrichName from .enrich_rfm import EnrichRFM from .enrich_superfan import EnrichSuperfan from .enrich_unique_gender import EnrichUniqueGender from .enrichment import Enrichment log = logging.getLogger() ENRICHMENTS = class_dict(__name__, Enrichment) def enrich_query(endpoint, req, event): collection_id = req["collection_id"] schema = req.get("workspace_schema") or req.get("alliance_schema") user_id = req["user_id"] enrichment_config = req.get("enrichment_config", {}) enr = endpoint["enrichment"](collection_id, schema, user_id, **enrichment_config) result = enr.do_everything() return result def segment_query(endpoint, req, event): collection_id = req["collection_id"] schema = req.get("workspace_schema") or req.get("alliance_schema") http_query(endpoint, req, event) process_status = get_enrichment_progress(schema, collection_id) appsync.send_process_status(schema, collection_id, process_status) def change_collection_status_endpoint(endpoint, req, event): collection_id = req["collection_id"] schema = req.get("alliance_schema") or req.get("workspace_schema") user_id = req["user_id"] status = req["status"] change_collection_status(schema, user_id, collection_id, status) def make_enrichment_endpoints(): return { f"enrich/{enr_name}": {"enrichment": enr_class, "query": enrich_query} for enr_name, enr_class in ENRICHMENTS.items() } def make_post_enrichment_endpoints(): return { "post_enrich/create_segments": { "url": DP_URL + "create_segments", "method": "POST", "query": segment_query, }, "post_enrich/calculate_analytics_cache": { "url": DP_URL + "calculate_analytics_cache", "method": "POST", "query": http_query, }, "post_enrich/change_status": {"query": change_collection_status_endpoint}, }