import json from smelog.factory import BoundLogger from juno_filters_messages.clients import service from juno_filters_messages.constants import DATA_FIELDS, TTL def parse_events(lambda_event: dict) -> list: result = [] for record in lambda_event.get("Records", []): body = json.loads(record["body"]) event = json.loads(body.get("Message", {})) if event: result.append(event) return result def handler(logger: BoundLogger, lambda_event: dict): """Lambda job. Args: logger: Logger instance. lambda_event: Lambda event obj. Returns: JSON serializable response """ event_list = parse_events(lambda_event) if event_list: settings_list = service.user_data.get_settings() else: settings_list = [] for event in event_list: event_id = event["id"] sent_message_list = [ (i["account_id"], i["meta"]["filter_id"]) for i in service.user_data.get_messages(event_id) ] logger.debug(f"{event_id}: {len(sent_message_list)} messages found") new_message_list = [] for settings in settings_list: user_id, account_id, settings_id, filter_map = ( settings["user_id"], settings["account_id"], settings["settings"]["id"], settings["settings"]["data"].get("juno_filters", {}).get("filters", {}), ) for filter_id, filter_data in filter_map.items(): if not filter_data.get("is_subscribed") or (account_id, filter_id) in sent_message_list: continue new_message_list.append( { "ttl": TTL, "event_id": event_id, "account_id": account_id, "meta": { "subject": "juno_digest", "views": ["juno_email"], "user_id": user_id, "settings_id": settings_id, "filter_id": filter_id, }, "data": { key: value for key, value in filter_data.items() if key in DATA_FIELDS }, } ) logger.debug(f"{len(new_message_list)} messages to send") if new_message_list: service.user_data.post_messages(data=new_message_list)