import os import sys import traceback from datetime import datetime from typing import Dict, List, Set import boto3 import redis import sentry_sdk from apollo_notifications.decorators import loop, time_logger from apollo_notifications.main_db import session_scope from apollo_notifications.playlists.serializers import SpotifyPlaylistUpdatePushSchema from apollo_notifications.playlists.utils import filter_playlist_update_event from apollo_notifications.push_client.client import PushClient from apollo_notifications.push_client.data_classes import PlaylistUpdatePushData, Push from apollo_notifications.push_client.exceptions import PushClientError from apollo_notifications.redis import ConstantKeyCache, lock from apollo_notifications.user_data.client import UserDataClient from apollo_notifications.utils import dump_datetime from ddtrace import patch_all, tracer from sentry_sdk.utils import BadDsn from sqlalchemy.exc import SQLAlchemyError from client import SpotifyPlaylistClient from config import config from logger import logger if os.environ.get('DATADOG_SERVICE_NAME'): patch_all() else: tracer.enabled = False try: if config.ENVIRONMENT != 'local': sentry_sdk.init(dsn=config.SENTRY_DSN, environment=config.ENVIRONMENT) except BadDsn: pass redis_client = redis.StrictRedis(config.REDIS_HOST, config.REDIS_PORT) def process_starred_playlists_updates_by_market( count: int, market: str, playlists: List[object], existing_messages: Set[str], notifications_client: SpotifyPlaylistClient, push_client: PushClient, users_map: Dict[str, int] = None, ): schema = SpotifyPlaylistUpdatePushSchema( title=config.TITLE, topic=config.TOPIC, vendor=config.VENDOR, message_template=config.MESSAGE_TEMPLATE, market=market, push_cls=Push, push_data_cls=PlaylistUpdatePushData, date=None, # for compatibility, ignored, real date is taken from dumped item users_map=users_map, ) messages = notifications_client.get_push_messages( query=playlists, markets=[market], topics=[config.TOPIC], push_schema=schema, existing_messages=existing_messages, filter_function=filter_playlist_update_event ) push_client.process_messages(messages, track_id_getter=lambda m: m.data.playlist_id) logger.info(f"2.{count}.1 Created {len(messages)} starred playlists update messages for {market} market.") @loop(logger, config) def job( notifications_client: SpotifyPlaylistClient, push_client: PushClient, user_data_client: UserDataClient, blacklist_cache: ConstantKeyCache ): # '_job' is separated from 'job' to allow its testing without wrapping decorators return _job( notifications_client=notifications_client, push_client=push_client, user_data_client=user_data_client, blacklist_cache=blacklist_cache ) def _job( notifications_client: SpotifyPlaylistClient, push_client: PushClient, user_data_client: UserDataClient, blacklist_cache: ConstantKeyCache ): _, users_map = user_data_client.parse_v1_mobile_settings(config.VENDOR, with_account_id=True) users_list = list(users_map.keys()) logger.info( f"1.0. Got {len(users_list)} users to send starred playlist additions notifications for: {users_list}.") starred_playlists_to_users_map, _ = user_data_client.parse_starred_playlists(config.VENDOR, users_list) if \ users_list else {} if not starred_playlists_to_users_map: logger.info("No starred playlist were found. Finished.") return with session_scope(): blacklisted_ids = set(blacklist_cache(notifications_client.get_blacklisted_ids)()) logger.info(f"1.1. Got {len(blacklisted_ids)} blacklisted playlists: {blacklisted_ids}.") last_update_datetime_str, is_not_found = notifications_client.get_last_update_datetime() current_update_datetime_str = dump_datetime(datetime.utcnow()) existing_messages = notifications_client.get_existing_push_messages( topics=[config.TOPIC], min_push_date=last_update_datetime_str, max_push_date=current_update_datetime_str) logger.info(f"1.4. Got {len(existing_messages)} existing push messages from {last_update_datetime_str} " f"to {current_update_datetime_str} for {config.TOPIC} - {config.VENDOR}") for p_id in blacklisted_ids: starred_playlists_to_users_map.pop(p_id, None) logger.info( f"1.2. Got {len(starred_playlists_to_users_map)} starred playlists to create notifications " f"for: {starred_playlists_to_users_map.keys()}.") market_to_playlists_map = notifications_client.get_market_to_updated_playlists_map( min_datetime=last_update_datetime_str, max_datetime=current_update_datetime_str, starred_playlists_to_users_map=starred_playlists_to_users_map ) for i, market in enumerate(market_to_playlists_map.keys()): try: with session_scope() as session: push_client.set_db_session(session) process_starred_playlists_updates_by_market( count=i, market=market, playlists=market_to_playlists_map[market], existing_messages=existing_messages, notifications_client=notifications_client, push_client=push_client, users_map=users_map, ) except (SQLAlchemyError, PushClientError) as ex: sentry_sdk.capture_exception(ex) logger.error(traceback.format_exception(*sys.exc_info()) + traceback.format_stack()) with session_scope(): notifications_client.set_last_update_datetime(current_update_datetime_str, create=is_not_found) @time_logger(logger, config.APP_NAME) @lock(logger, config, redis_client, sentry=sentry_sdk) def handler(): """Save push messages to db and send them to queue.""" sqs_client = boto3.client('sqs') user_data_client = UserDataClient(config, logger) notifications_client = SpotifyPlaylistClient(config) push_client = PushClient(logger=logger, sqs_client=sqs_client, config=config) blacklist_cache = ConstantKeyCache(redis_client, config.BLACKLIST_KEY, config.BLACKLIST_TTL) job(notifications_client, push_client, user_data_client, blacklist_cache) if __name__ == "__main__": handler()