from typing import Any, Dict, List, Mapping from analytics.constants import cache from analytics.constants import date as date_constants from analytics.queries.format import format_row from analytics.queries.product import ( ProductBulkGrowthPeriods, ProductDownloadsTimeSeries, ProductStreamsTimeSeries, ) from analytics.utils.cache import cache_in_redis from analytics.utils.store_availability import is_apple_spotify_data_in_sync # TODO implement remaining aggregations from analytics.utils.streams import ( breakdown_by_sos, calc_date_skip_rate, check_store_ids_for_sos_detailed, get_streams_sos_detailed_columns, get_streams_sos_detailed_columns_from_stream_sources, ) TOTAL_TIMESERIES_TYPE = "PRODUCT_STREAMS" TOTAL_SUMMARY_TYPE = "TOTAL" AGGREGATION_FIELDS_TIMESERIES = { "PRODUCT_STREAMS_BY_STORE": "store_id", "PRODUCT_STREAMS_BY_TRACK": "tracks.isrc", "PRODUCT_STREAMS_BY_COUNTRY": "country_code", "PRODUCT_DOWNLOADS_BY_COUNTRY": "country_code", "PRODUCT_DOWNLOADS_BY_STORE": "store_id", "TRACK_DOWNLOADS_BY_TRACK": "tracks.isrc", "TRACK_DOWNLOADS_BY_COUNTRY": "country_code", "TRACK_DOWNLOADS_BY_STORE": "store_id", } AGGREGATION_FIELDS_SUMMARY = { "TOTAL": "product_id", "STORE": "store_id", # TODO SOUND_RECORDING was "tracks.track_unique_id", check data and multi product "SOUND_RECORDING": "isrc", "COUNTRY": "country_code", # TODO check data is accurate when using this for id "SOS": "product_id", } # map of type -> aggregation type -> table PRODUCT_TABLES_TIMESERIES = { # Product streams "PRODUCT_STREAMS": { "total": "V_STREAMS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "country": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, "PRODUCT_STREAMS_BY_TRACK": { "total": "V_STREAMS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "country": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, "PRODUCT_STREAMS_BY_COUNTRY": { "total": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", "country": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, "PRODUCT_STREAMS_BY_STORE": { "total": "V_STREAMS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "country": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, "PRODUCT_STREAMS_BY_SOS": { "total": "V_STREAMS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "country": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, "PRODUCT_STREAMS_BY_SOS_DETAILED": { "total": "V_STREAMS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "country": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, # Product downloads "PRODUCT_DOWNLOADS": { "total": "DOWNLOADS_BY_PRODUCT_FEED_DISTRIBUTOR_DAILY", "country": "DOWNLOADS_BY_PRODUCT_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, "PRODUCT_DOWNLOADS_BY_COUNTRY": { "total": "DOWNLOADS_BY_PRODUCT_COUNTRY_FEED_DISTRIBUTOR_DAILY", "country": "DOWNLOADS_BY_PRODUCT_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, "PRODUCT_DOWNLOADS_BY_STORE": { "total": "DOWNLOADS_BY_PRODUCT_FEED_DISTRIBUTOR_DAILY", "country": "DOWNLOADS_BY_PRODUCT_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, # Product track downloads "TRACK_DOWNLOADS": { "total": "DOWNLOADS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "country": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, "TRACK_DOWNLOADS_BY_TRACK": { "total": "DOWNLOADS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "country": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, "TRACK_DOWNLOADS_BY_COUNTRY": { "total": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", "country": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, "TRACK_DOWNLOADS_BY_STORE": { "total": "DOWNLOADS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "country": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", }, } # map of type -> aggregation type -> dict of tables required PRODUCT_TABLES_SUMMARY = { # Product streams "TOTAL": { "total": { "streams_table": "V_STREAMS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "product_downloads_table": "DOWNLOADS_BY_PRODUCT_FEED_DISTRIBUTOR_DAILY", "product_track_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", # noqa }, "country": { "streams_table": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa "product_downloads_table": "DOWNLOADS_BY_PRODUCT_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa "product_track_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa }, }, "SOUND_RECORDING": { "total": { "streams_table": "V_STREAMS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "product_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", # noqa "product_track_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", # noqa }, "country": { "streams_table": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa "product_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa "product_track_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa }, }, "COUNTRY": { "total": { "streams_table": "STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", "product_downloads_table": "DOWNLOADS_BY_PRODUCT_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa "product_track_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa }, "country": { "streams_table": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa "product_downloads_table": "DOWNLOADS_BY_PRODUCT_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa "product_track_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa }, }, "STORE": { "total": { "streams_table": "V_STREAMS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "product_downloads_table": "DOWNLOADS_BY_PRODUCT_FEED_DISTRIBUTOR_DAILY", "product_track_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", # noqa }, "country": { "streams_table": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa "product_downloads_table": "DOWNLOADS_BY_PRODUCT_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa "product_track_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa }, }, "SOS": { "total": { "streams_table": "V_STREAMS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", "product_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", # noqa "product_track_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_FEED_DISTRIBUTOR_DAILY", # noqa }, "country": { "streams_table": "V_STREAMS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa "product_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa "product_track_downloads_table": "DOWNLOADS_BY_PRODUCT_TRACK_COUNTRY_FEED_DISTRIBUTOR_DAILY", # noqa }, }, } def get_product_table_for_timeseries(query_type: str, countries: List) -> str: key = "total" if len(countries) != 0: key = "country" table_options = PRODUCT_TABLES_TIMESERIES.get(query_type) if not table_options: raise NotImplementedError(f"Query {query_type} not supported") table = table_options.get(key) if not table: raise NotImplementedError(f"Query {query_type} not supported with type {key}") return table @cache_in_redis(ttl=cache.ONE_DAY) def get_product_timeseries( query_params: Mapping[str, Any], permissions: Mapping[str, Any], ) -> List[Dict]: if query_params.get("multi_product"): # check if we're requesting all time and if so, get the min release date if query_params.get("date_period") == date_constants.ALL_TIME_GRAPHQL: query_params["get_all_time_for_multi_product"] = True query_type = query_params.get("type", TOTAL_TIMESERIES_TYPE) aggregation_field = AGGREGATION_FIELDS_TIMESERIES.get(query_type) if aggregation_field: query_params["aggregation_field"] = aggregation_field countries = query_params.get("countries", []) query_params["product_table"] = get_product_table_for_timeseries( query_type, countries ) if query_type == "PRODUCT_STREAMS_BY_SOS_DETAILED": store_ids = check_store_ids_for_sos_detailed(query_params.get("store_ids", [])) stream_sources = query_params.get("stream_sources", []) if stream_sources: query_params[ "streams_sos_columns" ] = get_streams_sos_detailed_columns_from_stream_sources( store_ids, stream_sources ) else: query_params["streams_sos_columns"] = get_streams_sos_detailed_columns( store_ids ) if query_type.startswith("PRODUCT_STREAMS"): query = ProductStreamsTimeSeries({**query_params, **permissions}) else: query = ProductDownloadsTimeSeries({**query_params, **permissions}) time_series = query.execute() ts = [] # todo refactor for data_point in time_series: data_point = format_row(data_point) if ( data_point.get("streams_with_skips") or data_point.get("streams_with_skips") == 0 ): skip_rate = calc_date_skip_rate(data_point) del data_point["streams_with_skips"] data_point["skip_rate"] = skip_rate ts.append(data_point) if query_type == "PRODUCT_STREAMS_BY_SOS": ids = query_params.get("ids") timeseries = [] for item in ts: timeseries.extend(breakdown_by_sos(item, sources=ids)) ts = timeseries if query_type == "PRODUCT_STREAMS_BY_SOS_DETAILED": timeseries = [] for item in ts: timeseries.extend( breakdown_by_sos( item, sources=query_params["streams_sos_columns"], version=2, ) ) ts = timeseries return ts @cache_in_redis(ttl=cache.ONE_DAY) def get_product_growth_periods_bulk( query_params: Mapping[str, Any], permissions: Mapping[str, Any], ) -> List[Dict]: query_params["is_feed_data_available"] = is_apple_spotify_data_in_sync() query = ProductBulkGrowthPeriods({**query_params, **permissions}) return [format_row(data_point) for data_point in query.execute()]