"""Copy track playlist history data for today or up to a week ago.""" import os import sys import time import traceback from datetime import datetime, timedelta from exceptions import DspApiError, PlaylistUpdateError, WorkTimeExceeded from typing import Dict, List, Set, Tuple import sentry_sdk from ddtrace import patch_all, tracer from sentry_sdk.utils import BadDsn from sqlalchemy.exc import SQLAlchemyError import config import main_db from consts import APP_NAME, UpdateStatus from dsp_api_client import DspClient from logger import logger from redis_db import get_redis_lock 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 def get_tracks_to_add_map( client: DspClient, playlist_id: str, start_date: str, end_date: str, saved_isrc_set: Set[str] ) -> Tuple[Dict[str, List[str]], Set[str]]: """Get data for particular playlist from delphi using dsp-api and pack them to structures. :param client: dsp-api client instance to communicate with delphi api. :param playlist_id: id of playlist to update. :param start_date: string (ISO) start date for range to get streams for. :param end_date: string (ISO) end date for range to get streams for. :param saved_isrc_set: set of isrc for tracks which already presented for this playlist in db (if get same from delphi we should not update corresponding records, so them can be skipped). :return tuple of: - tracks_to_add_map - dict of isrc to list of 'added_datetime', 'track_id' for all tracks which should be inserted in track list of this playlist. - unchanged_isrc_set - set of tracks isrc which haven't changed presence in the playlist. """ streams_list = client.get_streams_list(playlist_id, start_date=start_date, end_date=end_date) tracks_enter_map = dict() unchanged_isrc_set = set() for item in streams_list: # if ordering (asc) by date will not work, should take min of all dates for this isrc instead first _isrc = item["isrc"].upper() if _isrc in saved_isrc_set: unchanged_isrc_set.add(_isrc) continue tracks_enter_map.setdefault(_isrc, [datetime.strptime(item["date"], "%Y-%m-%d")]) return tracks_enter_map, unchanged_isrc_set def process_playlist(client: DspClient, playlist_id: str, start_date: str, end_date: str) -> UpdateStatus: """Update particular playlists track list with tracks got from delphi. :param client: dsp-api client instance to communicate with delphi api. :param playlist_id: id of playlist to update. :param start_date: string (ISO) start date for range to get streams for. :param end_date: string (ISO) end date for range to get streams for. """ saved_tracks_isrc = set(main_db.get_saved_tracks_isrc_list(playlist_id)) isrc_enter_map, isrc_unchanged_set = get_tracks_to_add_map( client, playlist_id, start_date, end_date, saved_tracks_isrc ) isrc_exit_set = saved_tracks_isrc - isrc_unchanged_set if not isrc_enter_map and not isrc_exit_set: logger.info(f"Playlist {playlist_id} is unchanged.") return UpdateStatus.UNCHANGED unfilled_id_set = set(isrc_enter_map.keys()) for item in main_db.get_tracks_ids_query(list(isrc_enter_map.keys())): _isrc = item.isrc.upper() isrc_enter_map[_isrc].append(item.track_id) unfilled_id_set.remove(_isrc) if unfilled_id_set: logger.warning( f"Playlist {playlist_id}: impossible to get track id for " f"{len(unfilled_id_set)}/{len(isrc_enter_map.keys())} tracks: {unfilled_id_set}, skipped." ) for _isrc in unfilled_id_set: isrc_enter_map.pop(_isrc) main_db.delete_playlist_tracks(playlist_id, isrc_list=list(isrc_exit_set)) main_db.add_playlist_tracks(playlist_id, isrc_enter_map) logger.info(f"Playlist {playlist_id} was successfully updated.") return UpdateStatus.UPDATED def update_tracklist_data(): """Update personalized playlists track lists with streams data got from Delphi.""" lock = None try: if config.REDIS_LOCK_ENABLED: lock = get_redis_lock() if not lock.acquire(blocking=False): logger.info("Can not get lock") return start = time.time() dsp = DspClient(config.Dsp()) end_date = datetime.now() - timedelta(hours=config.END_DATE_HOURS_SHIFT) start_date_str = (end_date - timedelta(days=config.DATE_RANGE_DAYS_SHIFT)).strftime("%Y-%m-%d") end_date_str = end_date.strftime("%Y-%m-%d") playlist_id_list = main_db.get_personalized_playlists_list() errors_map = dict() updated_count = 0 for playlist_id in playlist_id_list: try: time.sleep(config.UPDATE_DELAY) update_status = process_playlist(dsp, playlist_id, start_date_str, end_date_str) updated_count += int(update_status) except (SQLAlchemyError, DspApiError): errors_map[playlist_id] = "".join( traceback.format_exception(*sys.exc_info()) + traceback.format_stack() ) continue elapsed_time = time.time() - start elapsed_time_str = time.strftime("%H:%M:%S", time.gmtime(elapsed_time)) log_str = ( f"{APP_NAME} finished in {elapsed_time_str}.\n" f"Of the {len(playlist_id_list)} personalized playlists found, " f"{len(playlist_id_list) - updated_count - len(errors_map)} were unchanged, \n" f"{updated_count} were updated, and " f"{len(errors_map)} caused errors when trying to update." ) logger.info(log_str) if errors_map: raise PlaylistUpdateError( log_str + "\n\n".join( [ f"Playlist {playlist_id} update got error\n{error_str}" for playlist_id, error_str in errors_map.items() ] ) ) if config.MAX_EXPECTED_WORK_TIME and elapsed_time >= config.MAX_EXPECTED_WORK_TIME: raise WorkTimeExceeded( f"{APP_NAME} successfully finished, but elapsed time {elapsed_time_str} " f"exceeded expected maximum {config.MAX_EXPECTED_WORK_TIME}" ) except Exception as ex: sentry_sdk.capture_exception(ex) raise finally: if config.REDIS_LOCK_ENABLED and lock: lock.release() if __name__ == "__main__": update_tracklist_data()