import logging from service.async_task_manager import io_task from service.async_task_manager.task import Task from service.tasks.analytics.modules import get_am from service.utils.data_model_utils import get_attribute_id logger = logging.getLogger(__name__) def handle_constant(conf, schema, **kwargs): return conf["value"] def handle_none(conf, schema, **kwargs): return None @io_task def handle_analytics(conf, schema, recalculate=False, **kwargs): conf["inputs"]["schema"] = schema if "field_name" in conf["inputs"]: conf["inputs"]["attribute_id"] = get_attribute_id( schema, conf["inputs"]["field_name"] ) am = get_am(conf) try: result = am.get_result(recalculate=recalculate) if result: out_map = conf["output_map"] am_map = am.get_output_map() if isinstance(out_map, str): val = result[0][am_map[out_map].idx] if val is not None: return val elif isinstance(out_map, dict): data_map = {d: am_map[s].idx for d, s in out_map.items()} return [ {dest: data[source] for dest, source in data_map.items()} for data in result ] except Exception as e: logger.exception(e) return conf["default"] PARAM_HANDLERS = { "constant": handle_constant, "analytics": handle_analytics, None: handle_none, } def get_param_handler(ptype: str): if ptype in PARAM_HANDLERS: return PARAM_HANDLERS[ptype] raise RuntimeError(f"Unknown filter parameter handler {ptype}") def finalize_param(param): if isinstance(param, Task): return param.wait_for_result() return param