import asyncio from aiohttp import web from aiohttp_apispec import docs, querystring_schema from apollo_utils.core.constants import ALL from collections import defaultdict from datetime import timedelta from server.client import services from server.client.utils import str_to_date from server.constants import BASE_API_PREFIX, DSP, cache_keys from server.constants.charts import ChartBreakdown, ChartType from server.constants.charts.delphi import TrackStateIncludes from server.constants.distributors import DISTRIBUTORS, IsSonyInclude from server.constants.market import Market from server.constants.tracks import TiktokAnalyticsBreakdowns, TrackPlaylistInclude, TracksGTPInclude from server.scenarios.distributors import get_tracks_distributors_map from server.schemas.tracks.charts import TracksCharts from server.schemas.tracks.gtp import TracksGTP from server.schemas.tracks.markets import TracksTopMarkets from server.schemas.tracks.new_music_friday import NMFTracks from server.schemas.tracks.playlists import TracksPlaylists from server.schemas.tracks.streams import TracksStreamsByDateShift from server.utils import sort_limit_dict from server.utils.pagination import multikeysort, paginate from server.utils.parallel import make_requests from server.utils.response import dump_response_schema router = web.RouteTableDef() @router.get(BASE_API_PREFIX + "/tracks/gtp/") @docs( tags=["gtp"], summary="Get gtp tracks.", ) @querystring_schema(TracksGTP.Request) @dump_response_schema(TracksGTP.Response, apply=True) async def get(request: web.Request) -> web.Response: params = request["querystring"] type_id, gtp_date, market, include = params["type_id"], params["gtp_date"], params["market"], params["include"] dsp_market = Market.convert_global(market, Market.WORLDWIDE) history_response = await services.apollo.get_gtp_tracks_history(type_id=type_id, date=gtp_date) gtp_date = history_response["gtp_date"] track_id_list = [r["id"] for r in history_response["tracks"]] track_isrc_list = [r["isrc"] for r in history_response["tracks"]] dates_include = TracksGTPInclude.DATES.value in include ( dsp_tracks, charts_isrc, tracks_weeks, streams_date, chart_date, hot_hits_date, tracks_re_entry, ) = await make_requests( ( (services.vendor.get_tracks, dict(dsp=DSP.SPOTIFY, ids=track_id_list, market=market)), ( services.dsp.get_tracks_in_chart, dict( isrc=track_isrc_list, chart_type=[ChartType.REGIONAL.value], chart_breakdown=[ChartBreakdown.DAILY.value], chart_country_code=[dsp_market], ), ), ( services.apollo.get_gtp_tracks_weeks, dict(isrc=track_isrc_list, date=gtp_date), {}, TracksGTPInclude.WEEKS.value in include, ), (services.dsp.get_streams_latest_date, {}, None, dates_include), ( services.dsp.get_charts_latest_date, dict( dsp=DSP.SPOTIFY, chart_type=ChartType.REGIONAL.value, chart_breakdown=ChartBreakdown.DAILY.value, country_code=dsp_market, ), None, dates_include, ), (services.apollo.get_gtp_hot_hits_latest_date, {}, None, dates_include), ( services.apollo.get_gtp_tracks_re_entry, dict(isrc=track_isrc_list, date=gtp_date), [], TracksGTPInclude.RE_ENTER.value in include, ), ) ) dsp_tracks = {t["id"]: t for t in dsp_tracks} for track_data in history_response["tracks"]: isrc = track_data["isrc"] track_data.update( { "data": dsp_tracks.get(track_data["id"]), "weeks": tracks_weeks.get(isrc, 0), "in_chart": isrc in charts_isrc, "is_re_enter": isrc in tracks_re_entry, } ) history_response.update( { "streams_date": streams_date, "chart_date": chart_date, "hot_hits_date": hot_hits_date, } ) return history_response @router.get(BASE_API_PREFIX + "/tracks/streams/by-date-shift/") @docs( tags=["streams"], summary="Get tracks streams by isrc, end date and period length in days.", ) @querystring_schema(TracksStreamsByDateShift.Request) @dump_response_schema(TracksStreamsByDateShift.Response(many=True)) async def get(request: web.Request) -> web.Response: params = request["querystring"] streams_date, date_shift, isrc_list, market = ( params["streams_date"], params["date_shift"], params["isrc_list"], params["market"], ) latest_date = await services.dsp.get_streams_latest_date() end_date = min(streams_date, latest_date) if latest_date else streams_date start_date = end_date - timedelta(date_shift) requested_dates = {start_date.isoformat(), end_date.isoformat()} streams_items = await services.dsp.get_streams_per_country( isrc=isrc_list, start_date=start_date, end_date=end_date, market=market, vendor=DSP.SPOTIFY.value, combine_isrc=False, streams_only=True, ) for isrc_item in streams_items: data = isrc_item.get("data", [{}])[0].get("data", []) for i in range(len(data) - 1, -1, -1): if data[i].get("date") not in requested_dates: data.pop(i) isrc_item["data"] = data return streams_items @router.get(BASE_API_PREFIX + "/tracks/playlists/") @docs( tags=["playlists"], summary="Get track playlists.", ) @querystring_schema(TracksPlaylists.Request) @dump_response_schema(TracksPlaylists.Response) @paginate(cache_keys.TRACK_PLAYLISTS) async def get(request: web.Request) -> web.Response: params = request["querystring"] isrc_list, dsp, include, markets_list, streams_markets_list = ( params["isrc_list"], params["dsp"], params["include"], params["markets_list"], params["streams_markets_list"], ) tasks = [] include_track_streams = TrackPlaylistInclude.TRACK_STREAMS.value in include if include_track_streams: tasks.append(services.dsp.get_streams_latest_date()) include.remove(TrackPlaylistInclude.TRACK_STREAMS.value) tasks.append( services.apollo.get_track_playlists( dsp, isrc=isrc_list, market=markets_list, search=params["search"], recent_adds_only=params["recent_adds_only"], include=include, **{"image_size": params["image_size"]} if dsp == DSP.APPLE else {"streams_market": streams_markets_list}, ) ) responses = await asyncio.gather(*tasks) result = responses[-1]["items"] if include_track_streams: end_date = responses[0] start_date = end_date - timedelta(days=6) streams_items = await services.dsp.get_streams( isrc=isrc_list, start_date=start_date, end_date=end_date, country_code=streams_markets_list if streams_markets_list else [Market.WORLDWIDE], dsp=dsp.value, subset="playlists", ) streams_mapping = defaultdict(int) for item in streams_items: streams_mapping[item["playlist_id"].replace(f"{dsp.value}_", "")] += item.get("streams", 0) for item in result: item["track_streams"] = streams_mapping[item["id"]] result = multikeysort(result, params["order_by"]) return result @router.get(BASE_API_PREFIX + "/tracks/nmf/") @docs(tags=["tracks", "nmf"], summary="Get NMF tracks.") @querystring_schema(NMFTracks.Request) @dump_response_schema(NMFTracks.Response) async def get(request: web.Request) -> web.Response: params = request["querystring"] market, distributors = params.pop("market"), params.pop("distributors", []) result = await services.apollo.get_nmf_tracklist(**params) if not result.get("items"): return result dsp_track_id_list = [f"{DSP.SPOTIFY.value}_{i['trackId']}" for i in result["items"]] track_id_to_distributor, extra_data = await get_tracks_distributors_map( country_code=market, dsp_track_id=dsp_track_id_list, distributors=set(distributors) | {DISTRIBUTORS.SME}, include=(IsSonyInclude.TRACK_LIST,), remove_result_dsp_prefix=True, ) track_id_to_meta = {t["id"]: t for t in extra_data.get(IsSonyInclude.TRACK_LIST, [])} for item in result["items"]: track_id = item["trackId"] track_data = track_id_to_meta.get(track_id, {}) distributor = track_id_to_distributor.get(track_id) album_image = None for image_data in track_data.get("album", {}).get("images", []): if not album_image or album_image.get("width", 0) > image_data.get("width", 0): album_image = image_data item.update( { "trackName": track_data.get("name"), "sonyRelease": None, "distributed_by": distributor, "isSony": distributor == DISTRIBUTORS.SME.value, "artists": [{"id": t["id"], "name": t["name"]} for t in track_data.get("artists", [])], "albumImageUrl": album_image.get("url") if album_image else None, } ) return result @router.get(BASE_API_PREFIX + "/tracks/markets/top/") @docs(tags=["tracks", "markets"], summary="Get tracks top markets.") @querystring_schema(TracksTopMarkets.Request) @dump_response_schema(TracksTopMarkets.Response) async def get(request: web.Request) -> web.Response: params = request["querystring"] start_date, end_date, isrc_list, dsp_list, combine, include_worldwide, metrics_list, limit = ( params["start_date"], params["end_date"], params["isrc_list"], params["dsp_list"], params["combine"], params["include_worldwide"], params["metrics_list"], params["limit"], ) include_tiktok = DSP.TIKTOK.value in dsp_list and metrics_list tasks = [] if include_tiktok: tasks.append( services.dsp.get_tiktok_tracks_analytics( start_date=start_date, end_date=end_date, isrc=isrc_list, metrics=metrics_list, breakdowns=TiktokAnalyticsBreakdowns.COUNTRY_TOTALS.value, ) ) dsp_list.remove(DSP.TIKTOK.value) if dsp_list: tasks.append(services.dsp.get_streams(start_date=start_date, end_date=end_date, isrc=isrc_list, dsp=dsp_list)) responses = await asyncio.gather(*tasks) market_metrics = defaultdict(lambda: defaultdict(int)) if dsp_list: for item in responses[-1]: market = item["country_code"] if include_worldwide or market != Market.WORLDWIDE: market_metrics[ALL if combine else item["dsp"]][market] += item.get("streams", 0) if include_tiktok: tiktok_metrics = {metric: {} for metric in metrics_list} for market, market_values in responses[0]["breakdowns"][TiktokAnalyticsBreakdowns.COUNTRY_TOTALS.value].items(): if include_worldwide or market != Market.WORLDWIDE: if combine: market_metrics[ALL][market] += sum(market_values.values()) else: for metric in metrics_list: tiktok_metrics[metric][market] = market_values.get(metric, 0) if not combine: for metric in metrics_list: tiktok_metrics[metric] = [ {"market": market, metric: value} for market, value in sort_limit_dict(tiktok_metrics[metric], limit) ] if len(metrics_list) == 1: tiktok_metrics = tiktok_metrics[metrics_list[0]] market_metrics[DSP.TIKTOK.value] = tiktok_metrics result = { dsp: ( [{"market": market, "streams": streams} for market, streams in sort_limit_dict(markets_streams, limit)] if dsp != DSP.TIKTOK.value else markets_streams ) for dsp, markets_streams in market_metrics.items() } return web.json_response(result) @router.get(BASE_API_PREFIX + "/tracks/charts/") @docs(tags=["tracks", "charts"], summary="Get track's charts.") @querystring_schema(TracksCharts.Request) @dump_response_schema(TracksCharts.Response, apply=True) async def get(request: web.Request) -> web.Response: params = request["querystring"] isrc, dsp, chart_breakdown, chart_date = ( params["isrc"], params["dsp"], params.get("chart_breakdown"), params.pop("chart_date"), ) chart_types = ( dict(chart_type=ChartType.REGIONAL.value, chart_breakdown=chart_breakdown.value) if dsp.value == DSP.SPOTIFY.value else {} ) all_chart_list, track_chart_list, lifetime_chart_list, country_code_list = await asyncio.gather( *[ services.dsp.get_dsp_charts(dsp, type=ChartType.REGIONAL.value, breakdown=chart_breakdown.value), services.dsp.get_tracks_charts( dsp, isrc=isrc, date=chart_date, track_state_includes=[ TrackStateIncludes.CHART_META.value, TrackStateIncludes.METRICS.value, TrackStateIncludes.LIFETIME_METRICS.value, ], metrics_dimension="isrc", **chart_types, ), services.dsp.get_tracks_charts_lifetime( dsp, isrc=isrc, track_state_includes=[TrackStateIncludes.CHART_META.value, TrackStateIncludes.LIFETIME_METRICS.value], metrics_dimension="isrc", **chart_types, ), services.apollo.get_markets(region_type=f"charts_{dsp.value}", extended=True, mapping_only=False), ] ) track_chart_map = {i["chart_meta"]["chart_id"]: i for i in track_chart_list} lifetime_chart_map = {i["chart_meta"]["chart_id"]: i for i in lifetime_chart_list} country_code_map = {item["code"]: item for item in country_code_list} result = [] for item in all_chart_list: chart_id, country_code = item["chart_id"], item["country_code"] # current market if chart_id in track_chart_map: result_item = track_chart_map[chart_id] # removed market - entered before the chosen date but not in chart right now elif chart_id in lifetime_chart_map and ( not chart_date or str_to_date(lifetime_chart_map[chart_id]["lifetime_metrics"]["earliest_position_date"]) < chart_date ): result_item = lifetime_chart_map[chart_id] # not present market else: result_item = {"chart_meta": item} result.append(result_item) apollo_country_code = Market.GL if country_code == Market.WORLDWIDE else country_code if apollo_country_code in country_code_map: country_code_item = country_code_map[apollo_country_code] result_item["chart_meta"].update( {"name": country_code_item["name"], "region": country_code_item.get("region")} ) return {"items": result}