import asyncio from aiohttp import web from aiohttp_apispec import docs, json_schema, querystring_schema, response_schema from apollo_utils.core.constants.dsp import DSP from apollo_utils.core.constants.market import Market from apollo_utils.service.clients.aiohttp.utils.response import dump_response_schema from collections import defaultdict from datetime import datetime, timedelta, timezone from typing import Dict from server.constants.delphi.streams.group_by import StreamsGroupBy from server.constants.delphi.streams.include import StreamsInclude from server.constants.delphi.streams.subset import StreamsSubset from server.constants.delphi.tiktok import TiktokTrackAnalyticsBreakdown from server.constants.delphi.videos.expand_to import VideosExpandTo from server.constants.delphi.videos.group_by import VideosGroupBy from server.legacy.analytics import deserializers from server.legacy.analytics import utils as analytics_utils from server.legacy.analytics import view_utils from server.legacy.analytics.constants import ( DETAILED_STREAMS_MONITORING_SOURCE_V2, TIKTOK_METRICS_LIST, TWO_WEEKS_DAYS, YoutubeVideosOrder, ) from server.legacy.analytics.schemas import ( BulkTracksStreamsSchemaV1, DetailedTrackStreamsGraphViewV2, MetricsMarketsWeekValuesSchema, MetricsTrendsSchema, PlaylistsStreamsGraphSchema, ShazamCharts, ShazamChartsPositions, ShazamCities, ShazamCountries, SinglePlaylistDateRangeStreamsSchema, StreamsMonitoringSchemaV1, TikTokBaseRequestSchema, TikTokCountryCodeSchema, TikTokCountryCodeTopSchema, TikTokGraphResponseSchema, TiktokResponseInsightSchema, TiktokShortSummarySchema, TikTokTrendsResponseSchema, TotalStreamsCountSchemaV1, TrackInPlaylistStreamsSchema, TrackMarketsListGraphSchemaV1, TrackMarketsStreamsListSchemaV1, TrackPerformance, TrackPlaylistsStreams, TracksStationsStreamsSummarySchema, TrackStreamsGraph, TrackStreamsGraphViewV1, TrackTopPlaylistsStreamsSchema, TrackVideos, YoutubeDemographics, YoutubeShortSummarySchema, YoutubeTopMarketsSchema, ) from server.legacy.analytics.utils import has_tiktok_none_data, get_amazon_playlist_worldwide_streams_graph from server.legacy.analytics.youtube_utils import ( calculate_trends, define_biggest_source, define_demographics_popularity, fill_missing_dates, map_and_group_analytics_data_by_date, map_and_group_demographics_percent_data, ) from server.legacy.core.constants import ( ALL_VENDORS, APPLE, DELPHI_GLOBAL_MARKET, DELPHI_ZZ_MARKET, SPOTIFY, TIKTOK, YOUTUBE, ) from server.legacy.core.utils import form_chunk_lists, format_delphi_playlist_id, multikeysort from server.legacy.delphi.constants import SortBy, VideosOnly @docs( tags=["streams"], summary="Returns total num of streams for specific market and global one for all 3 vendors.", description="Calculate total num of streams for track y it ISRC value in global and specific market.", ) @querystring_schema(TotalStreamsCountSchemaV1.RequestSchema) @response_schema(TotalStreamsCountSchemaV1.ResponseSchema, description="") async def total_streams_count_view(request: web.Request) -> web.Response: return await view_utils.get_total_streams_count( request, ALL_VENDORS, first_stream_dates=False, global_field="global_streams", market_field="market_streams", ) @querystring_schema(deserializers.DailyWeeklyStreamsDeserializer) async def daily_and_weekly_based_streams_data_view( request: web.Request, ) -> web.Response: return await view_utils.get_daily_and_weekly_based_streams_data(request, ALL_VENDORS, with_dates=True) @querystring_schema(deserializers.StreamsDatesDeserializer) async def streams_dates_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] include_last_dates = data["include_last_dates"] tasks = [api.get_first_stream_dates(isrc_list=data["isrc_list"], dsp_list=ALL_VENDORS)] if include_last_dates: tasks.append(api.get_latest_stream_dates()) responses = await asyncio.gather(*tasks) first_stream_date = defaultdict(dict) for item in responses[0]: first_stream_date[item["dsp"]][item["isrc"]] = item["date"] result = {"first_stream_date": first_stream_date} if include_last_dates: result["last_stream_date"] = {i["dsp"]: i["updated_date"] for i in responses[1]} return web.json_response(result) @querystring_schema(deserializers.SimpleWeekStreamsDeserializer) async def weeks_based_streams_data_view(request: web.Request) -> web.Response: return await view_utils.get_daily_and_weekly_based_streams_data( request, ALL_VENDORS, with_dates=True, days_and_chart_weeks=False ) @querystring_schema(deserializers.DateRangeStreamsDeserializer) async def date_range_streams_view(request: web.Request) -> web.Response: """Total streams of num in specific range of dates and for specific or global market.""" return await view_utils.get_date_range_streams(request, ALL_VENDORS, streams_field="streams") @docs( tags=["streams", "tracks", "bulk"], summary="Returns dict with list of ISRCs and streams for them.", description="Calculate streams for each day in specific date range across all 3 vendors for a list of tracks.", ) @querystring_schema(BulkTracksStreamsSchemaV1.RequestSchema) @response_schema(BulkTracksStreamsSchemaV1.ResponseSchema) async def tracks_range_streams_view(request: web.Request) -> web.Response: return await view_utils.get_tracks_streams(request) @docs( tags=["streams", "graph", "performance"], summary="Returns coordinates to build performance streams graph in specific date range for specific market and " + "global one for all 3 vendors.", description="Calculate streams for each day in specific date range across all 3 vendors.", ) @querystring_schema(TrackStreamsGraphViewV1.RequestSchema) @response_schema(TrackStreamsGraphViewV1.ResponseSchema) async def vendors_streams_graph_view(request: web.Request) -> web.Response: return await view_utils.get_stream_graph(request, ALL_VENDORS, show_total=True, add_market=True, default_none=True) @querystring_schema(StreamsMonitoringSchemaV1.RequestSchema) @response_schema(StreamsMonitoringSchemaV1.ResponseSchema, description="") async def streams_monitoring_data_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] response = await api.get_streams( start_date=data["start_date"], end_date=data["end_date"], isrc_list=[data["isrc"]], dsp_list=ALL_VENDORS, country_code_list=[data["market"]], include=[StreamsInclude.ENGAGEMENT], ) results = {v: None for v in ALL_VENDORS} for item in response: vendor = item["dsp"] streams_info = item[f"{item['dsp']}_streams_info"] results[vendor] = streams_info.get("engagement") or streams_info.get("listener_engagement") return web.json_response(results) @docs( tags=["streams", "graph", "performance"], summary="Returns coordinates to build detailed performance streams graph in specific date range for specific " + "market and global one for selected vendor.", description="Calculate all streams sources for each day in specific date range for selected vendor.", ) @querystring_schema(DetailedTrackStreamsGraphViewV2.RequestSchema) @response_schema(DetailedTrackStreamsGraphViewV2.ResponseSchema) async def vendor_detailed_streams_graph_view_v2(request: web.Request) -> web.Response: return await view_utils.get_stream_graph( request, ALL_VENDORS, with_sources=True, wrap_single_vendor=True, add_market=True, source_mapping=DETAILED_STREAMS_MONITORING_SOURCE_V2, ) @docs( tags=["streams", "graph", "markets"], summary="Returns list of coordinates for selected markets.", description="Returns list of coordinates for selected markets.", ) @querystring_schema(TrackMarketsListGraphSchemaV1.RequestSchema) @response_schema(TrackMarketsListGraphSchemaV1.ResponseSchema) async def markets_streams_graph_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] date_range = analytics_utils.DateRange(data["start_date"], data["end_date"]) markets = data["markets"] response = await api.get_streams( start_date=date_range.start, end_date=date_range.end, isrc_list=[data["isrc"]], dsp_list=ALL_VENDORS, country_code_list=markets, group_by=[StreamsGroupBy.DATE], ) result = { "markets": analytics_utils.process_stream_graph_data_per_market( response, date_range.start, date_range.end, ALL_VENDORS, markets, data["per_vendors"], ), "dates": date_range.to_dict(isoformat=True), } return web.json_response(result) @docs( tags=["streams", "track", "markets"], summary="Returns list of ordered by streams markets for track, possible to retrieve all or limited num.", description="Returns list of ordered by streams markets for track.", ) @querystring_schema(TrackMarketsStreamsListSchemaV1.RequestSchema) @response_schema(TrackMarketsStreamsListSchemaV1.ResponseSchema) async def track_markets_streams_list_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] limit = data["limit"] current_range = analytics_utils.DateRange(start=data["start_date"], end=data["end_date"]) prev_range = analytics_utils.get_previous_date_range(current_range) response = await api.get_streams( start_date=prev_range.start, end_date=current_range.end, isrc_list=[data["isrc"]], dsp_list=data["vendors"], group_by=[StreamsGroupBy.DATE], ) markets_data = analytics_utils.collect_markets_data(response, current_range, prev_range) markets_data = sorted(markets_data, key=lambda m: m["current_streams"], reverse=True) result = { "markets": markets_data[:limit] if limit else markets_data, "dates": { "current_range": current_range.to_dict(isoformat=True), "prev_range": prev_range.to_dict(isoformat=True), }, } return web.json_response(result) @docs( tags=["streams", "tracks", "stations"], summary="Get a list of stations with streams combined for all ISRC for chosen market plus global or all.", description="Returns a list of tracks stations.", ) @querystring_schema(TracksStationsStreamsSummarySchema.RequestSchema) @response_schema(TracksStationsStreamsSummarySchema.ResponseSchema) async def tracks_stations_streams_summary_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] isrc_list = data["isrc_list"] market = data["market"] days_list = data["days_list"] min_streams = data["min_streams"] compact_result = data["compact"] vendors = [APPLE] end_date = await api.get_latest_stream_date_vendor(vendor=APPLE) tasks = [ api.get_streams( isrc_list=[isrc], dsp_list=vendors, country_code_list=[DELPHI_GLOBAL_MARKET, market] if market else None, start_date=(end_date - timedelta(days=days_count - 1)), end_date=end_date, subset=StreamsSubset.PLAYLISTS, ) for days_count in days_list for isrc in isrc_list ] responses = await asyncio.gather(*tasks) playlist_mask = f"{APPLE}_pl." results = defaultdict(lambda: defaultdict(lambda: {days: 0 for days in days_list})) for i, days_count in enumerate(days_list): for j, isrc in enumerate(isrc_list): data = responses[i * len(isrc_list) + j] for item in data: if not item["playlist_id"].startswith(playlist_mask): results[item["playlist_id"]][item["country_code"]][days_count] += item.get("streams", 0) max_days = max(days_list) return web.json_response( { "items": [ { "id": playlist_id, "total": list(markets_data[DELPHI_GLOBAL_MARKET].values()), "markets": [ [market] + list(streams.values()) if compact_result else {"market": market, "streams": list(streams.values())} for market, streams in markets_data.items() if market != DELPHI_GLOBAL_MARKET ], } for playlist_id, markets_data in results.items() if min_streams is None or markets_data[DELPHI_GLOBAL_MARKET][max_days] >= min_streams ], "last_date": end_date.isoformat(), } ) @querystring_schema(TrackInPlaylistStreamsSchema.RequestSchema) @response_schema(TrackInPlaylistStreamsSchema.ResponseSchema) async def tracks_in_playlist_streams_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] vendor, playlist_id = data["vendor"], data["playlist_id"] current_range = analytics_utils.DateRange(start=data["start_date"], end=data["end_date"]) prev_range = analytics_utils.get_previous_date_range(current_range) tasks = [ api.get_streams( isrc_list=data["isrc_list"], dsp_list=[vendor], country_code_list=data["markets"], start_date=date_range.start, end_date=date_range.end, subset=StreamsSubset.PLAYLISTS, playlist_id_list=[format_delphi_playlist_id(vendor, playlist_id)], ) for date_range in [current_range, prev_range] ] responses = await asyncio.gather(*tasks) items_list = analytics_utils.process_track_in_playlist_streams_data(responses) result = { "items": items_list, "playlist_id": playlist_id, "dates": { "current_range": current_range.to_dict(isoformat=True), "prev_range": prev_range.to_dict(isoformat=True), }, } return web.json_response(result) @docs( tags=["streams", "video", "youtube"], summary="Get summary statistic with trends for selected/last day, day before, current week and previous week.", description="Returns a dict with YT statistics for several periods.", ) @querystring_schema(YoutubeShortSummarySchema.RequestSchema) @response_schema(YoutubeShortSummarySchema.ResponseSchema) async def youtube_short_summary_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] latest_date = data["latest_date"] if not latest_date: latest_date = await api.get_latest_stream_date_vendor(YOUTUBE) date_ranges = analytics_utils.get_short_summary_intervals(latest_date) _, _, last_week, previous_week = date_ranges response = await api.get_video_analytics( start_date=previous_week.start, end_date=last_week.end, isrc=[data["isrc"]], country_code_list=data.get("country_code_list"), content_type=data["content_type_list"], expand_to=[VideosExpandTo.RELATED_ISRCS], group_by=[VideosGroupBy.VIDEO_ID, VideosGroupBy.DATE], only=data["only"], ) results = analytics_utils.calc_youtube_short_summary(response, date_ranges) return web.json_response(analytics_utils.format_short_summary(results, date_ranges)) @docs( tags=["streams", "track", "tiktok"], summary="Get summary statistics with trends for selected/last day, day before, current week and previous week.", description="Returns a dict with Tiktok statistics for several periods.", ) @querystring_schema(TiktokShortSummarySchema.RequestSchema) @response_schema(TiktokShortSummarySchema.ResponseSchema) async def tiktok_short_summary_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] latest_date = data["latest_date"] if not latest_date: latest_date = await api.get_latest_stream_date_vendor(TIKTOK) date_ranges = analytics_utils.get_short_summary_intervals(latest_date) _, _, last_week, previous_week = date_ranges response = await api.get_tiktok_tracks_analytics( start_date=previous_week.start, end_date=last_week.end, isrc_list=[data["isrc"]], country_code_list=data.get("country_code_list"), content_type_list=data["content_type_list"], breakdowns_list=[TiktokTrackAnalyticsBreakdown.CONTENT_TYPE_COUNTRY], metrics_list=data["only"], ) response = response["breakdowns"][TiktokTrackAnalyticsBreakdown.CONTENT_TYPE_COUNTRY] results = analytics_utils.calc_tiktok_short_summary(response) return web.json_response(analytics_utils.format_short_summary(results, date_ranges)) @docs( tags=["streams", "video", "youtube"], summary="Get top markets by views.", description="Returns top markets by views for selected dates range and videos.", ) @querystring_schema(YoutubeTopMarketsSchema.RequestSchema) @response_schema(YoutubeTopMarketsSchema.ResponseSchema) async def youtube_top_markets_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] skip_markets = None if data["include_worldwide"] else (DELPHI_GLOBAL_MARKET, DELPHI_ZZ_MARKET) limit, keep_zero_streams, include, sort_by = ( data["limit"], data["keep_zero_streams"], data.get("include"), data.get("sort_by"), ) response = await api.get_video_analytics( start_date=data["start_date"], end_date=data["end_date"], video_id=data["video_id_list"], only=[VideosOnly.VIEWS], ) views_by_market = defaultdict(int) total_views = 0 for item in response: country_code = item["dimensions"]["country_code"] if not skip_markets or country_code not in skip_markets: views_cnt = item["metrics"]["views"] views_by_market[country_code] += views_cnt total_views += views_cnt views_by_market = [ {"market": i[0], "views": i[1]} for i in views_by_market.items() if i[1] > 0 or keep_zero_streams ] if "percentage" in include: for item in views_by_market: if item["views"] > 0: item["percentage"] = round(100 * item["views"] / total_views) # default for mobile ["-views", "market"] result = multikeysort(views_by_market, sort_by)[:limit] return web.json_response({"items": result}) @docs( tags=["streams", "track", "playlists"], summary="Get track playlists streams for selected markets in specified date range.", description="Returns a list with track playlists streams for selected markets.", ) @json_schema(TrackPlaylistsStreams.Request) @response_schema(TrackPlaylistsStreams.Response) async def tracks_playlist_streams_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["json"] vendor, isrc = data["vendor"], data.get("isrc") playlists_markets = data["items"] current_range = analytics_utils.DateRange(start=data["start_date"], end=data["end_date"]) prev_range = analytics_utils.get_previous_date_range(current_range) playlists_ids, markets = [], set() chunk_size = data["chunk_size"] for _id, _markets in playlists_markets.items(): playlists_ids.append(format_delphi_playlist_id(vendor, _id)) markets.update(set(_markets)) current_range_data, prev_range_data = [], [] for chunk_list in form_chunk_lists(playlists_ids, chunk_size): _markets = list(markets) if markets else None tasks = [ api.get_streams( isrc_list=[isrc] if isrc else None, dsp_list=[vendor], country_code_list=_markets, start_date=date_range.start, end_date=date_range.end, subset=StreamsSubset.PLAYLISTS, playlist_id_list=chunk_list, ) for date_range in [current_range, prev_range] ] _current_range_chunk, _prev_range_chunk = await asyncio.gather(*tasks) current_range_data.extend(_current_range_chunk) prev_range_data.extend(_prev_range_chunk) items_list = analytics_utils.process_track_playlists_streams_data( (current_range_data, prev_range_data), vendor, playlists_markets ) result = { "items": items_list, "dates": { "current_range": current_range.to_dict(isoformat=True), "prev_range": prev_range.to_dict(isoformat=True), }, } return web.json_response(TrackPlaylistsStreams.Response().dump(result)) @docs( tags=["streams", "track", "playlists", "top"], summary="Get top list of track playlists in specified date range.", description="Returns a streams ordered list with track playlists streams data for selected markets.", ) @json_schema(TrackTopPlaylistsStreamsSchema.RequestSchema) @response_schema(TrackTopPlaylistsStreamsSchema.ResponseSchema) async def top_track_playlists_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["json"] vendor, isrc = data["vendor"], data["isrc"] chunk_size = data["chunk_size"] markets, playlists_ids = data["markets"], data["playlists_ids"] playlists_ids = [format_delphi_playlist_id(vendor, pk) for pk in playlists_ids] if playlists_ids else [None] current_range = analytics_utils.DateRange(start=data["start_date"], end=data["end_date"]) prev_range = analytics_utils.get_previous_date_range(current_range) current_range_data, prev_range_data = [], [] for chunk_list in form_chunk_lists(playlists_ids, chunk_size): _playlists = chunk_list if next(iter(chunk_list), None) else None tasks = [ api.get_streams( isrc_list=[isrc], dsp_list=[vendor], country_code_list=markets, start_date=date_range.start, end_date=date_range.end, subset=StreamsSubset.PLAYLISTS, sort_by=SortBy.STREAMS, playlist_id_list=_playlists, ) for date_range in [current_range, prev_range] ] _current_range_chunk, _prev_range_chunk = await asyncio.gather(*tasks) current_range_data.extend(_current_range_chunk) prev_range_data.extend(_prev_range_chunk) items_dict = analytics_utils.process_track_top_playlists_streams_data([current_range_data, prev_range_data], vendor) items = list(items_dict.values()) result = { "count": len(items), "items": items, "dates": { "current_range": current_range.to_dict(isoformat=True), "prev_range": prev_range.to_dict(isoformat=True), }, } return web.json_response(TrackTopPlaylistsStreamsSchema.ResponseSchema().dump(result)) @docs( tags=["streams", "playlists", "graph"], summary="Playlists streams graph endpoint.", description="Returns graph coordinates for requested playlists in specified markets.", ) @querystring_schema(PlaylistsStreamsGraphSchema.RequestSchema) @response_schema(PlaylistsStreamsGraphSchema.ResponseSchema) async def playlists_streams_graph_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] vendor = data["vendor"] markets, playlists_ids = data["markets"], data["playlists_ids"] vendor_playlists_ids = [format_delphi_playlist_id(vendor, pk) for pk in playlists_ids] date_range = analytics_utils.DateRange(start=data["start_date"], end=data["end_date"]) group_market_as, group_playlist_id_as = data["group_market_as"], data["group_playlist_id_as"] response = await api.get_streams( dsp_list=[vendor], country_code_list=markets, start_date=date_range.start, end_date=date_range.end, subset=StreamsSubset.PLAYLISTS, group_by=[StreamsGroupBy.DATE], playlist_id_list=vendor_playlists_ids, ) if DSP.AMAZON.value in vendor and group_market_as: response = get_amazon_playlist_worldwide_streams_graph( response, playlists_ids, group_playlist_id_as, group_market_as ) result = { "items": analytics_utils.process_playlists_streams_for_graph(response, date_range, vendor), "dates": date_range.to_dict(isoformat=True), } return web.json_response(PlaylistsStreamsGraphSchema.ResponseSchema().dump(result)) @docs( tags=["streams", "playlist", "range"], summary="Playlist date ranges streams endpoint.", description="Returns playlist current and previous date range streams data.", ) @querystring_schema(SinglePlaylistDateRangeStreamsSchema.RequestSchema) @response_schema(SinglePlaylistDateRangeStreamsSchema.ResponseSchema) async def playlist_date_range_streams_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] vendor, market = data["vendor"], data["market"] current_range = analytics_utils.DateRange(start=data["start_date"], end=data["end_date"]) prev_range = analytics_utils.get_previous_date_range(current_range) vendor_playlist_id = format_delphi_playlist_id(vendor, data["playlist_id"]) tasks = [ api.get_streams( dsp_list=[vendor], country_code_list=[market], start_date=date_range.start, end_date=date_range.end, subset=StreamsSubset.PLAYLISTS, playlist_id_list=[vendor_playlist_id], ) for date_range in (current_range, prev_range) ] current_range_data, prev_range_data = await asyncio.gather(*tasks) current_data = next(iter(current_range_data), {}) prev_data = next(iter(prev_range_data), {}) result = { "market": market, "current_range": current_data.get("streams"), "prev_range": prev_data.get("streams"), "dates": { "current_range": current_range.to_dict(isoformat=True), "prev_range": prev_range.to_dict(isoformat=True), }, } return web.json_response(SinglePlaylistDateRangeStreamsSchema.ResponseSchema().dump(result)) @docs( tags=["streams", "track", "tiktok", "spotify"], summary="Get day and week trends for Spotify and Tiktok.", description="Returns trends stats per each metric.", ) @querystring_schema(MetricsTrendsSchema.RequestSchema) @response_schema(MetricsTrendsSchema.ResponseSchema) async def metrics_trends_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] isrc_list, metric_list, latest_date, market = ( data["isrc_list"], data["metric_list"], data["latest_date"], data["market"], ) if not latest_date: latest_date = await api.get_min_latest_date([SPOTIFY, TIKTOK]) period = 7 week_start_date = latest_date - timedelta(days=(period - 1)) start_date = week_start_date - timedelta(days=period) metrics_data = await analytics_utils.get_metrics_trends_data( api, isrc_list, market, metric_list, start_date, latest_date ) result = analytics_utils.calc_metrics_trends(metrics_data, metric_list, week_start_date, latest_date) return web.json_response(MetricsTrendsSchema.ResponseSchema().dump({"items": result})) @docs( tags=["streams", "track", "tiktok", "spotify"], summary="Get 7 days Spotify/Tiktok metrics values sums for worldwide and selected or top market.", description="Returns metrics week sums for two markets.", ) @querystring_schema(MetricsMarketsWeekValuesSchema.RequestSchema) @response_schema(MetricsMarketsWeekValuesSchema.ResponseSchema) async def metrics_markets_week_values_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] data = request["querystring"] latest_date, market = ( data["latest_date"], data["market"], ) if not latest_date: latest_date = await api.get_min_latest_date([SPOTIFY, TIKTOK]) start_date = latest_date - timedelta(days=6) metrics_data = await analytics_utils.get_metrics_markets_week_values_data( api, data["isrc_list"], market, data["metric_list"], start_date, latest_date ) result = analytics_utils.calc_metrics_markets_week_values(metrics_data, market, start_date, latest_date) return web.json_response(MetricsMarketsWeekValuesSchema.ResponseSchema().dump({"items": result})) @docs( tags=["shazam", "country_codes"], summary="Get shazam country code list.", description="Returns shazam country codes.", ) @response_schema(ShazamCountries.Response) async def shazam_countries_view(request: web.Request) -> web.Response: api = request.app["delphi_java_api"] result = await api.get_shazam_cities(sort_by="country_code") result = list(set(i["country_code"] for i in result)) return web.json_response(ShazamCountries.Response().dump({"items": result})) @docs( tags=["shazam", "cities"], summary="Get shazam cities for specific market.", description="Returns shazam cities.", ) @querystring_schema(ShazamCities.Request) @response_schema(ShazamCities.Response) async def shazam_cities_view(request: web.Request) -> web.Response: api, data = request.app["delphi_java_api"], request["querystring"] market = data["market"] result = await api.get_shazam_cities(country_code=market, sort_by="city_name") return web.json_response(ShazamCities.Response().dump({"market": market, "items": result})) @docs( tags=["shazam", "charts"], summary="Get shazam charts data.", description="Returns shazam chart metrics.", ) @querystring_schema(ShazamCharts.Request) @response_schema(ShazamCharts.Response) async def shazam_charts_view(request: web.Request) -> web.Response: api, data = request.app["delphi_java_api"], request["querystring"] result = await view_utils.get_shazam_charts_positions_and_summary( api, analytics_utils.get_shazam_chart_id(data["market"], data["city"]), data["chart_date"], data["is_sony"], ) return web.json_response(ShazamCharts.Response().dump({"items": result})) @docs( tags=["shazam", "charts"], summary="Get shazam charts positions for dates range.", description="Returns shazam charts positions list.", ) @querystring_schema(ShazamChartsPositions.Request) @response_schema(ShazamChartsPositions.Response) async def shazam_charts_positions_view(request: web.Request) -> web.Response: result = await view_utils.get_shazam_charts_position(request) context = { "min_position": result["position"]["peak"], "max_position": result["position"]["min"], } return web.json_response(ShazamChartsPositions.Response(context=context).dump(result)) @docs(tags=["youtube"], summary="Youtube track videos list.") @querystring_schema(TrackVideos.Request) @dump_response_schema(TrackVideos.Response) async def youtube_get_videos_list(request: web.Request): params = request["querystring"] isrc, views = params["isrc"], params["views"] limit, offset, sort_by = params.get("limit"), params["offset"], params["sort_by"] delphi_java_api = request.app["delphi_java_api"] items = await delphi_java_api.get_videos( isrc_list=[isrc], views=views, content_type=["partner_uploaded", "premium_ugc"], expand_to=["related_isrcs"], ) if YoutubeVideosOrder.VIDEO_TITLE_DESC.value not in sort_by and YoutubeVideosOrder.VIDEO_TITLE.value not in sort_by: sort_by.append("video_title") items = multikeysort(items, sort_by) if limit: items = items[offset : offset + limit] return {"items": items} @docs(tags=["youtube"], summary="Youtube track performance data.") @querystring_schema(TrackPerformance.Request) @dump_response_schema(TrackPerformance.Response) async def youtube_get_performance(request: web.Request): params = request["querystring"] video_id, latest_date = (params["video_id_list"], params["latest_date"]) country_code_list = params.get("country_code_list", [Market.WORLDWIDE]) start_date = latest_date - timedelta(days=(TWO_WEEKS_DAYS * 2 - 1)) delphi_java_api = request.app["delphi_java_api"] performance_data = await delphi_java_api.get_video_analytics( group_by=["date", "video_id"], video_id=video_id, start_date=start_date, end_date=latest_date, country_code_list=country_code_list, only=["traffic_sources", "engagement", "watch_times"], ) performance_data = map_and_group_analytics_data_by_date(performance_data) performance_data = fill_missing_dates(performance_data, latest_date, TWO_WEEKS_DAYS * 2) trends = calculate_trends(performance_data) # We only need 4 weeks for trends performance_data = performance_data[:TWO_WEEKS_DAYS] biggest_traffic_source = define_biggest_source(performance_data) daily_metrics = [ { "date": i["date"], "timestampt": int(datetime.strptime(i["date"], "%Y-%m-%d").replace(tzinfo=timezone.utc).timestamp()), **i["metrics"], } for i in performance_data ] result = { "biggest_traffic_source": biggest_traffic_source, "daily_metrics": daily_metrics, "trends": trends, } return result @docs(tags=["youtube"], summary="Youtube track streams graph.") @querystring_schema(TrackStreamsGraph.Request) @dump_response_schema(TrackStreamsGraph.Response) async def youtube_get_streams_graph(request: web.Request): params = request["querystring"] video_id, latest_date = (params["video_id_list"], params["latest_date"]) country_code_list = params.get("country_code_list", [Market.WORLDWIDE]) start_date = latest_date - timedelta(days=(TWO_WEEKS_DAYS - 1)) delphi_java_api = request.app["delphi_java_api"] performance_data = await delphi_java_api.get_video_analytics( group_by=["date", "video_id"], video_id=video_id, start_date=start_date, end_date=latest_date, country_code_list=country_code_list, ) performance_data = map_and_group_analytics_data_by_date(performance_data) performance_data = fill_missing_dates(performance_data, latest_date, TWO_WEEKS_DAYS) result = [] for row in performance_data: x, y = None, None metrics = row["metrics"] views = metrics.get("views") if views: y = {"total": views, **metrics["traffic_source_types"]} x = int(datetime.strptime(row["date"], "%Y-%m-%d").replace(tzinfo=timezone.utc).timestamp()) result.append({"x": x, "y": y}) result = sorted(result, key=lambda d: d["x"]) return {"coordinates": result} @docs( tags=["youtube"], summary="Groups of users, among whom the video is the most popular.", ) @querystring_schema(YoutubeDemographics.Request) @dump_response_schema(YoutubeDemographics.Response) async def youtube_get_demographics(request: web.Request): params = request["querystring"] video_id, start_date, end_date = ( params["video_id_list"], params["start_date"], params["end_date"], ) country_code_list = params.get("country_code_list", [Market.WORLDWIDE]) delphi_java_api = request.app["delphi_java_api"] demographics_data = await delphi_java_api.get_video_analytics( group_by=["video_id"], video_id=video_id, start_date=start_date, end_date=end_date, country_code_list=country_code_list, content_type=["partner_uploaded", "premium_ugc"], only=["demographics"], ) result = { "gender_most_popular": None, "gender_percentage": None, "age_band_most_popular": None, "age_band_percentage": None, } if demographics_data: data = map_and_group_demographics_percent_data(demographics_data) if data["views"] != 0 and sum(data["gender_views_percentages"].values()) > 0: result = define_demographics_popularity(data) return result @docs( tags=["tiktok"], summary="Get day and 2 weeks trends for Tiktok.", description="Returns trends stats for tiktok per 2 metrics.", ) @querystring_schema(TikTokBaseRequestSchema) @dump_response_schema(TikTokTrendsResponseSchema) async def tiktok_trends_view(request: web.Request) -> Dict: api = request.app["delphi_java_api"] data = request["querystring"] isrc_list, latest_date = (data["isrc_list"], data["latest_date"]) country_code_list = data.get("country_code_list", [Market.WORLDWIDE]) if not latest_date: latest_date = await api.get_min_latest_date([TIKTOK]) week_start_date = latest_date - timedelta(days=(TWO_WEEKS_DAYS - 1)) start_date = week_start_date - timedelta(days=TWO_WEEKS_DAYS) tiktok_data = await api.get_tiktok_tracks_analytics( start_date=start_date, end_date=latest_date, isrc_list=isrc_list, country_code_list=country_code_list, breakdowns_list=[TiktokTrackAnalyticsBreakdown.DAILY], metrics_list=TIKTOK_METRICS_LIST, use_args_dates=True, ) metrics_data = tiktok_data["breakdowns"][TiktokTrackAnalyticsBreakdown.DAILY] # AG-10835 Decided not to calculate weekly trend only if data for the period is absent, # if we have at least data for one day from 14 days period the trend is calculated # that is why we need to pass period=1 to calc_metrics_trends result = analytics_utils.calc_metrics_trends(metrics_data, TIKTOK_METRICS_LIST, week_start_date, latest_date, 1) final_result = {} for row in result: key = row.pop("metric") final_result[key] = row final_result[key]["latest"]["trend"] = final_result[key]["latest"].pop("trend_one_day_ago") return final_result @docs( tags=["tiktok"], summary="Get 2 week insights for Tiktok.", description="Returns insights for tiktok per 2 metrics.", ) @querystring_schema(TikTokBaseRequestSchema) @dump_response_schema(TiktokResponseInsightSchema) async def tiktok_insights_view(request: web.Request) -> Dict: api = request.app["delphi_java_api"] data = request["querystring"] isrc_list, latest_date = (data["isrc_list"], data["latest_date"]) country_code_list = data.get("country_code_list", [Market.WORLDWIDE]) if not latest_date: latest_date = await api.get_min_latest_date([TIKTOK]) start_date = latest_date - timedelta(days=(TWO_WEEKS_DAYS - 1)) tiktok_data = await api.get_tiktok_tracks_analytics( start_date=start_date, end_date=latest_date, isrc_list=isrc_list, country_code_list=country_code_list, breakdowns_list=[TiktokTrackAnalyticsBreakdown.DAILY], metrics_list=TIKTOK_METRICS_LIST, use_args_dates=True, ) if not has_tiktok_none_data(tiktok_data): return { "creations": None, "views": None, "peak_date": None, } metrics_data = tiktok_data["breakdowns"][TiktokTrackAnalyticsBreakdown.DAILY] result = analytics_utils.calc_tiktok_insight(metrics_data, start_date) return result @docs( tags=["tiktok"], summary="Get 2 week graph for Tiktok.", description="Returns graph for tiktok per 2 metrics.", ) @querystring_schema(TikTokBaseRequestSchema) @dump_response_schema(TikTokGraphResponseSchema) async def tiktok_graph_view(request: web.Request) -> Dict: api = request.app["delphi_java_api"] data = request["querystring"] isrc_list, latest_date = (data["isrc_list"], data["latest_date"]) country_code_list = data.get("country_code_list", [Market.WORLDWIDE]) if not latest_date: latest_date = await api.get_min_latest_date([TIKTOK]) start_date = latest_date - timedelta(days=(TWO_WEEKS_DAYS - 1)) tiktok_data = await api.get_tiktok_tracks_analytics( start_date=start_date, end_date=latest_date, isrc_list=isrc_list, country_code_list=country_code_list, breakdowns_list=[ TiktokTrackAnalyticsBreakdown.DAILY, TiktokTrackAnalyticsBreakdown.TOTALS, ], metrics_list=TIKTOK_METRICS_LIST, use_args_dates=True, ) if not has_tiktok_none_data(tiktok_data): return { "creations": None, "video_views": None, } result = analytics_utils.build_tiktok_graph(tiktok_data, start_date) return result @docs( tags=["tiktok"], summary="Get 2 week top markets for Tiktok.", description="Returns 2 week top markets for tiktok per 2 metrics.", ) @querystring_schema(TikTokCountryCodeTopSchema.RequestSchema) @dump_response_schema(TikTokCountryCodeTopSchema.ResponseSchema) async def tiktok_top_markets_view(request: web.Request) -> Dict: api = request.app["delphi_java_api"] data = request["querystring"] isrc_list, latest_date, limit = ( data["isrc_list"], data["latest_date"], data["limit"], ) if not latest_date: latest_date = await api.get_min_latest_date([TIKTOK]) start_date = latest_date - timedelta(days=(TWO_WEEKS_DAYS - 1)) tiktok_data = await api.get_tiktok_tracks_analytics( start_date=start_date, end_date=latest_date, isrc_list=isrc_list, breakdowns_list=[TiktokTrackAnalyticsBreakdown.COUNTRY_TOTALS], metrics_list=TIKTOK_METRICS_LIST, use_args_dates=True, ) market_data = tiktok_data["breakdowns"][TiktokTrackAnalyticsBreakdown.COUNTRY_TOTALS] if not market_data: return { "creations": None, "video_views": None, } result = analytics_utils.calc_tiktok_top(market_data, limit) result["latest_date"] = latest_date return result @docs( tags=["tiktok"], summary="Get 2 week list of markets for Tiktok.", description="Returns 2 week list of markets for tiktok.", ) @querystring_schema(TikTokCountryCodeSchema.RequestSchema) @dump_response_schema(TikTokCountryCodeSchema.ResponseSchema) async def tiktok_markets_view(request: web.Request) -> Dict: api = request.app["delphi_java_api"] data = request["querystring"] isrc_list, latest_date = ( data["isrc_list"], data["latest_date"], ) if not latest_date: latest_date = await api.get_min_latest_date([TIKTOK]) start_date = latest_date - timedelta(days=(TWO_WEEKS_DAYS - 1)) tiktok_data = await api.get_tiktok_tracks_analytics( start_date=start_date, end_date=latest_date, isrc_list=isrc_list, breakdowns_list=[TiktokTrackAnalyticsBreakdown.COUNTRY_TOTALS], metrics_list=TIKTOK_METRICS_LIST, use_args_dates=True, ) market_data = tiktok_data["breakdowns"][TiktokTrackAnalyticsBreakdown.COUNTRY_TOTALS] country_codes = [key for key, val in market_data.items() if val["creations"] > 0 or val["video_views"] > 0] return {"country_codes": country_codes}