import logging import time from typing import Dict, List import boto3 import sentry_sdk from apollo_main_db.apollo.models import StarredContent, ApolloRecentSearch from apollo_main_db.apple.models import AppleMusicSong from apollo_main_db.spotify.models import SpotifyTrack2 from dsp import DigitalServiceProvider from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration from sentry_sdk.utils import BadDsn from sqlalchemy import func from config import Config from db import session_scope logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) SPOTIFY_TRACK_URI = 'spotify:track:' TRACK_SEARCH_TYPE = 2 config = Config() try: sentry_sdk.init(dsn=config.SENTRY_DSN, integrations=[AwsLambdaIntegration()]) except BadDsn: pass def get_user_ids() -> List[str]: """Get user_ids from s3.""" s3_client = boto3.client('s3') obj = s3_client.get_object(Bucket=config.BUCKET_NAME, Key=config.BUCKET_KEY) value = obj['Body'].read().decode('utf-8') return value.splitlines() def get_starred_tracks(session, user_ids: List[str]) -> Dict[str, str]: """Get starred tracks by user_ids""" query = session.query( func.REPLACE(StarredContent.uri, SPOTIFY_TRACK_URI, '').label('track_id'), (func.IF(AppleMusicSong.isrc.is_(None), SpotifyTrack2.isrc, AppleMusicSong.isrc)).label('isrc'), ).outerjoin( AppleMusicSong, AppleMusicSong.id == StarredContent.uri ).outerjoin( SpotifyTrack2, SpotifyTrack2.id == func.REPLACE(StarredContent.uri, SPOTIFY_TRACK_URI, '') ).filter( StarredContent.user_id.in_(user_ids), ).group_by(StarredContent.uri) return {q.track_id: q.isrc for q in query} def get_recent_search_tracks(session, user_ids: List[str]) -> Dict[str, str]: """Get recent tracks from search by user_ids""" tracks = {} for user_id in user_ids: query = session.query( func.REPLACE(ApolloRecentSearch.uri, SPOTIFY_TRACK_URI, '') ).filter( ApolloRecentSearch.user_id == user_id, ApolloRecentSearch.search_type == TRACK_SEARCH_TYPE ).group_by( ApolloRecentSearch.uri ).order_by( ApolloRecentSearch.created_at.desc() ).limit(10) for track, in query: tracks[track] = None query = session.query( SpotifyTrack2.id.label('track_id'), SpotifyTrack2.isrc.label('isrc'), ).filter( SpotifyTrack2.id.in_(tracks), ) return {**tracks, **{q.track_id: q.isrc for q in query}} def main() -> None: """Script entry point.""" logger.info('Start') start_time = time.time() dsp = DigitalServiceProvider(config=config, logger=logger) user_ids = get_user_ids() logger.info(f'Getting tracks for users: {user_ids}') with session_scope(config) as session: starred_tracks = get_starred_tracks(session, user_ids) recent_tracks = get_recent_search_tracks(session, user_ids) tracks = {**starred_tracks, **recent_tracks} isrc_list = [] tracks_without_isrc = [track_id for track_id, isrc in tracks.items() if not isrc] if tracks_without_isrc: isrc_list = dsp.get_isrc_for_spotify_tracks(tracks_without_isrc) tracks_with_isrc = filter(lambda x: x, tracks.values()) isrc_list = set(isrc_list) | set(tracks_with_isrc) for isrc in isrc_list: dsp.get_total_streams_by_isrc(isrc) time.sleep(1) logger.info(f'{isrc}: done') logger.info(f'Finished {time.time() - start_time} sec') if __name__ == '__main__': main()