import json from datetime import date, timedelta from logging import Logger from typing import Dict, List from apollo_utils.core.constants.market import Market from juno_email_messages import amplitude_events from juno_email_messages.clients import service from juno_email_messages.constants import DATE_FORMAT from juno_email_messages.render import generate_amplitude_event, generate_email 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 get_events(lambda_event: dict, logger: Logger) -> List[dict]: event_list = parse_events(lambda_event) event_id_list = [event["id"] for event in event_list if event.get("id")] if event_id_list: logger.debug(f"{len(event_id_list)} lambda events: {event_id_list}") else: logger.debug(f"No lambda events: {lambda_event}") return [] event_list = service.user_data.get_messages(event_id_list) if event_list: logger.debug(f"Got {len(event_list)} events from user-data-api") else: logger.debug("No events in user-data-api") return [] return event_list def prepare_juno_request(juno_date: date, event_data: dict) -> dict: country_code_list = event_data.get("country_code") return { **( {"country_code": Market.WORLDWIDE if Market.WORLDWIDE in country_code_list else country_code_list} if country_code_list else {} ), **{ key: event_data[key] for key in ("isrc_country_code", "percent_change") if event_data.get(key) is not None }, **{ key.replace("days", "date"): ( juno_date - timedelta(days=event_data[key] + (1 if key == "max_product_sale_days" else 0)) ) for key in ("min_product_sale_days", "max_product_sale_days") if event_data.get(key) is not None }, } def send_emails( event_list: List[dict], juno_date: date, code_name_mapping: Dict[str, str], id_email_mapping: Dict[str, str], logger: Logger, ): email_list = [] message_list = [] for event in event_list: event_meta = event["meta"] user_id, account_id, filter_id, event_data = ( event_meta["user_id"], event["account_id"], event_meta["filter_id"], event["data"] ) email_to = id_email_mapping.get(user_id) if not email_to: logger.debug(f"{account_id}/{filter_id}: Missing email {event_meta}") continue juno_request_data = prepare_juno_request(juno_date, event_data) juno_track_list = service.gate.get_juno_tracks(**juno_request_data) if not juno_track_list: logger.debug(f"{account_id}/{filter_id}: No juno data {juno_request_data} for {event_data}") continue logger.debug(f"{account_id}/{filter_id}: {event_data} email {email_to} tracks {len(juno_track_list)}") email_body = generate_email(filter_id, event_data, code_name_mapping, juno_date, juno_track_list) email_list.append( { "to": [email_to], "subject": f"[JUNO] Trending Tracks {juno_date.strftime(DATE_FORMAT)} - {event_data['name']}", "body": email_body, } ) message_list.append(generate_amplitude_event(user_id, filter_id, event_data)) result = service.notifications.send_email(email_list) logger.debug(f"{len(email_list)} sent: {result}") amplitude_events.send_amplitude_events(message_list, logger) def handler(lambda_event: dict, logger: Logger): """Lambda job. Args: logger: Logger instance. lambda_event: Lambda event obj. Returns: JSON serializable response """ service.logger = logger event_list = get_events(lambda_event, logger) if not event_list: return service.notifications.set_token(service.atlas.get_access_token()) user_id_list = list(set(event["meta"]["user_id"] for event in event_list)) id_email_mapping = service.apollo.get_users_emails(user_id_list) code_name_mapping = service.apollo.get_markets() juno_date = service.dsp.get_juno_date() logger.debug(f"juno date {juno_date}") send_emails(event_list, juno_date, code_name_mapping, id_email_mapping, logger) def check_email_body(): filter_id = "123456" event_data = { "name": "TestFilter1", "country_code": ["us"], "isrc_country_code": ["us", "ca"], "percent_change": 20, "min_product_sale_days": 365, "max_product_sale_days": 1, } juno_date = service.dsp.get_juno_date() code_name_mapping = service.apollo.get_markets() juno_request_data = prepare_juno_request(juno_date, event_data) juno_track_list = service.gate.get_juno_tracks(**juno_request_data) email_body = generate_email(filter_id, event_data, code_name_mapping, juno_date, juno_track_list) with open("test_email.html", "w") as f: f.write(email_body)