import asyncio from collections import defaultdict from datetime import date, datetime, timedelta, timezone from typing import Any, Dict, Iterable, List, Tuple, Union from apollo_utils.core.constants.dsp import DSP from server.client.clients.delphi_client import DelphiClient from server.constants.delphi.streams.group_by import StreamsGroupBy from server.constants.delphi.tiktok import TiktokTrackAnalyticsBreakdown from server.constants.sorting import SortOrder from server.legacy.analytics import constants as analytics_constants from server.legacy.analytics.constants import TIKTOK_METRICS_LIST from server.legacy.core import constants, utils from server.legacy.core.utils import get_optional_min, multikeysort from server.legacy.delphi import constants as delphi_constants from server.utils.combine import get_combined_data class DateRange: """ class the store week range start and end dates. """ def __init__(self, start: date, end: date): self.start = start self.end = end def __contains__(self, item: Union[str, datetime, date]): if isinstance(item, str): item = datetime.strptime(item, "%Y-%m-%d").date() elif type(item) is datetime: item = item.date() return self.start <= item <= self.end @property def days(self): return (self.end - self.start).days + 1 @property def start_isoformat(self): return self.start.isoformat() @property def end_isoformat(self): return self.end.isoformat() def get_range_dates(self) -> List[str]: """ Returns all dates in week range is iso format. Returns: List[str]: List of iso formatted dates. """ result = [] day = self.start while day <= self.end: result.append(day.isoformat()) day += timedelta(days=1) return result def check_in_range(self, date_obj: date) -> bool: """ Check if date in range. Args: date_obj (date): Date object. Returns: bool: True if in range, otherwise - False. """ return self.start <= date_obj <= self.end def __str__(self): return f"{self.start_isoformat} - {self.end_isoformat}" def to_dict(self, isoformat: bool = False): """ Convert object to dict. Args: isoformat (bool): Convert date object to isoformat or not. Returns: dict: Result dictionary start and end range values. """ return { "start": self.start_isoformat if isoformat else self.start, "end": self.end_isoformat if isoformat else self.end, } async def get_last_date_value(api, data: dict, key: str = "last_date") -> datetime.date: """ Retrieve last date value from json or make additional request for it. Arguments: api: Consumer analytics api. data (dict): Json data. Keyword Arguments: key (str): Key for json field (default: {'last_date'}). Returns: datetime.date: Last date value. """ last_date = data.get(key) if not last_date: last_date = await api.get_latest_stream_date_vendor(vendor=constants.SPOTIFY) return last_date def process_vendor_streaming_data(vendor: str, data: list) -> dict: """Process streaming data based on it's vendor. Arguments: vendor: Vendor name. data: Delphi service streaming data. Returns: dict: Processed vendor data or empty dict. """ if vendor == constants.APPLE: return process_apple_amazon_streaming_data(data) elif vendor == constants.SPOTIFY: return process_spotify_streaming_data(data) elif vendor == constants.AMAZON: return process_apple_amazon_streaming_data(data) def process_spotify_streaming_data(data: list) -> dict: """Process Spotify streaming data. Arguments: data: Raw streaming data. Returns: dict: Processed Spotify streaming data. """ item = data[0] if data else {} source_info = item.get("spotify_streams_info", {}).get("source", {}) return { "total_streams_num": item.get("streams", 0), "active_streams_num": count_spotify_engagement_streams(source_info, analytics_constants.ACTIVE_STREAM_FIELDS), "playlist_streams_num": count_spotify_engagement_streams( source_info, analytics_constants.PLAYLIST_STREAM_FIELDS ), "non_playlist_streams_num": count_spotify_engagement_streams( source_info, analytics_constants.NON_PLAYLIST_STREAM_FIELDS ), } def count_spotify_engagement_streams(streaming_data: dict, fields_list: list) -> int: """Count Spotify passively or actively discovered num of streams. Arguments: streaming_data: Streaming data for a specific date. Returns: Num of streams. """ streams_counter = 0 for field in fields_list: streams_counter += streaming_data.get(field, 0) return streams_counter def process_apple_amazon_streaming_data(data: list) -> dict: """Process Amazon/Apple streaming data. Arguments: data (list): Raw streaming data. Returns: dict: Processed Amazon streaming data. """ return {"total_streams_num": int(data[0].get("streams", 0)) if data else 0} def combine_total_streams(data: Dict[str, Dict[str, int]]) -> Dict[str, int]: """Combine total streams. Args: data: Streams per vendor per ISRC. Returns: Total streams per vendor. """ result = defaultdict(int) for item in data.values(): for key, value in item.items(): result[key] += value return result def process_total_streams_count( items: List[dict], combine_isrc: bool = True ) -> Tuple[Dict[str, int], Dict[str, int]] or Tuple[Dict[str, Dict[str, int]], Dict[str, Dict[str, int]]]: """Count streams for global and given market. Arguments: items: List of data. combine_isrc: Return min first date per vendor, do not need ISRC in result. Returns: tuple: total and market streams. """ global_streams = defaultdict(dict) market_streams = defaultdict(dict) for item in items: if item.get("country_code") == constants.DELPHI_GLOBAL_MARKET: global_streams[item["isrc"]][item["dsp"]] = item.get("streams", 0) else: market_streams[item["isrc"]][item["dsp"]] = item.get("streams", 0) if combine_isrc: return combine_total_streams(global_streams), combine_total_streams(market_streams) return global_streams, market_streams def process_first_stream_dates( data: List[dict], combine_isrc: bool = True ) -> Dict[str, Dict[str, date]] or Dict[str, date]: """Get first stream date per vendor. Arguments: data: Streams data. combine_isrc: Return min first date per vendor, do not need ISRC in result. Returns: Vendor to first stream date mapping. """ result = defaultdict(dict) for item in data: result[item["isrc"]][item["dsp"]] = item["date"] if not result: return {} if combine_isrc: combined_result = {} for isrc_dates in result.values(): if combined_result: for vendor, dsp_date in isrc_dates.items(): combined_result[vendor] = get_optional_min(combined_result.get(vendor), dsp_date) else: combined_result = isrc_dates return combined_result return result def _generate_coordinates(data: dict, start_date: date, end_date: date, default_none: bool): """Generate set of coordinates for dates interval with streaming info. Args: data: Streaming data per date. start_date: Interval from. end_date: Interval to. default_none: Display empty Y value as None - if True, or as zero - if False. Returns: List of coordinates. """ coordinates_list = [] for ts in utils.dates_range(start_date, end_date, as_timestamp=True): y_value = data.get(ts, 0) if default_none and not y_value: y_value = None coordinates_list.append({"x": ts, "y": y_value}) return coordinates_list def _generate_coordinates_per_vendor( data_per_vendor: Dict[str, dict], dsp_list: List[str], start_date: date, end_date: date, default_none: bool, ) -> dict: """Generate set of coordinates for dates interval with streaming info for vendor list. Args: data_per_vendor: Streaming data per date and vendor. dsp_list: Vendor list. start_date: Interval from. end_date: Interval to. default_none: Display empty Y value as None - if True, or as zero - if False. Returns: Dict per vendor of coordinates lists. """ return { vendor: _generate_coordinates(data_per_vendor[vendor], start_date, end_date, default_none) for vendor in dsp_list } def process_stream_graph_data( items: List[dict], start_date: date or datetime, end_date: date or datetime, dsp_list: List[str], show_total: bool = False, per_vendors: bool = False, base_mapping: Dict[str, str] or List[str] = None, source_mapping: Dict[str, Dict[str, str] or List[str]] = None, default_none: bool = False, ) -> Dict[str, List[Dict[str, int]]]: """Collect data for stream graph. Arguments: items: List of responses with items. start_date: Date ISO8601-formatted. end_date: Date ISO8601-formatted. dsp_list: All vendors that we need to get streams for. show_total: Show sum and optionally per vendor. per_vendors: For total case show per vendor streams or not. base_mapping: Streams item base level mapping. source_mapping: Streams item source level mapping. default_none: Display empty Y value as None - if True, or as zero - if False. Returns: List of coordinates [{'x': 1571346000, 'y': 190612},...] where x is timestamp, y is streams number. """ results_total = defaultdict(int) results_per_vendor = defaultdict(dict) with_details = bool(base_mapping or source_mapping) for item in items: item_date = datetime.strptime(item["date"], "%Y-%m-%d") ts = utils.date_to_timestamp(item_date) item_vendor = item["dsp"] streams_count = item.get("streams", 0) results_per_vendor[item_vendor][ts] = ( process_streams_data( item, source_mapping.get(item_vendor), source_field=("selection_source_type" if item_vendor == constants.AMAZON else "source"), ) if with_details else streams_count ) results_total[ts] += streams_count if show_total: results = {"coordinates": _generate_coordinates(results_total, start_date, end_date, default_none)} if per_vendors: results["vendors"] = _generate_coordinates_per_vendor( results_per_vendor, dsp_list, start_date, end_date, default_none ) return results return _generate_coordinates_per_vendor(results_per_vendor, dsp_list, start_date, end_date, default_none) def process_stream_graph_data_per_market( items: List[dict], start_date: date or datetime, end_date: date or datetime, dsp_list: List[str], market_list: List[str], per_vendors: bool = False, ) -> List[dict]: """Collect data for stream graph per market. Arguments: start_date: Date ISO8601-formatted. end_date: Date ISO8601-formatted. dsp_list: All vendors that we need to get streams for. market_list: Markets list. items: List of responses with items. per_vendors: For total case show per vendor streams or not. Returns: List of coordinates [{'x': 1571346000, 'y': 190612},...] where x is timestamp, y is streams number. """ results_total = defaultdict(lambda: defaultdict(int)) results_per_vendor = defaultdict(lambda: defaultdict(dict)) for item in items: item_date = datetime.strptime(item["date"], "%Y-%m-%d") ts = utils.date_to_timestamp(item_date) streams_count = item.get("streams", 0) market = item["country_code"] results_per_vendor[market][item["dsp"]][ts] = streams_count results_total[market][ts] += streams_count results = [] for market in market_list: results_item = {"market": utils.convert_market(market, constants.GLOBAL_MARKET)} if per_vendors: results_item["vendors"] = _generate_coordinates_per_vendor( results_per_vendor[market], dsp_list, start_date, end_date, True ) results_item["coordinates"] = _generate_coordinates(results_total[market], start_date, end_date, True) results.append(results_item) return results def get_mapping(item: dict, mapping: Dict[str, str] or Iterable[str]) -> dict: """Convert dict into another according to mapping. Args: item: Original dict. mapping: Mapping or field list. Returns: Modified dict. """ result = {} for key, source in mapping.items() if isinstance(mapping, dict) else zip(mapping, mapping): _value = 0 if isinstance(source, str): _value = item.get(source, 0) elif isinstance(source, list): sub_sources_list = [item.get(i) for i in source if item.get(i) is not None] _value = sum(sub_sources_list) result[key] = _value return result def process_streams_data( item: dict, source_mapping: Dict[str, str] or List[str], source_field: str = "source", ) -> dict or None: """Process streams data item. Args: item: Streams data item. source_mapping: Mapping for source info level streams data items. source_field: Source field name. Returns: Converted streams item. """ source_info = item.get(f"{item['dsp']}_streams_info", {}).get(source_field, {}) mapped_sources = get_mapping(source_info, source_mapping) total_sum = sum(mapped_sources.values()) if not total_sum: return result = {"total": total_sum, **mapped_sources} return result def get_chart_weeks_ranges(last_date: datetime) -> Tuple[DateRange, DateRange]: """ Returns date ranges for two latest available chart weeks. Chart week - latest available 7 days from Thursday to Friday. Args: last_date (datetime.datetime): Last date charts are available in DB. Returns: Tuple[DateRange, DateRange]: Tuple with start and end dates for last 2 chart weeks. """ if last_date.weekday() >= 3: # if last date week day is Thursday or later current_week_start = last_date - timedelta(weeks=1, days=last_date.weekday() - 4) current_week_end = current_week_start + timedelta(days=6) else: current_week_end = last_date - timedelta(weeks=1, days=last_date.weekday() - 3) current_week_start = current_week_end - timedelta(days=6) current_week = DateRange(start=current_week_start, end=current_week_end) prev_week = DateRange( start=(current_week_start - timedelta(weeks=1)), end=(current_week_end - timedelta(weeks=1)), ) return current_week, prev_week def get_simple_weeks_ranges(last_date: datetime) -> Tuple[DateRange, DateRange]: """Return date ranges for two latest available weeks. Args: last_date (datetime.datetime): Last date charts are available in DB. Returns: Tuple[DateRange, DateRange]: Tuple with start and end dates for last 2 weeks. """ current_week_start = last_date - timedelta(days=6) current_week = DateRange(start=current_week_start, end=last_date) prev_week = DateRange( start=(current_week_start - timedelta(weeks=1)), end=(last_date - timedelta(weeks=1)), ) return current_week, prev_week def get_rolling_week(end_date: date = None) -> DateRange: """Returns rolling week (range from last Friday to end_date) dates range. Args: end_date (datetime.date): Rolling week last date value (default: date.today()). Returns: DateRange: Rolling week date range object. """ end_date = end_date or date.today() if end_date.weekday() < 4: rolling_week_start = end_date - timedelta(weeks=1, days=end_date.weekday() - 4) else: rolling_week_start = end_date - timedelta(days=end_date.weekday() - 4) return DateRange(start=rolling_week_start, end=end_date) def _generate_daily_weekly_streams_part( streams_mapping: Dict[str, Dict[str, Dict[str, int]]], field_name: str, vendor: str = None, all_vendors: bool = False, ) -> Dict[str, Dict[str, int]]: return { f"{field_name}_streams": { field: sum(streams[field_name].values()) if all_vendors else streams[field_name].get(vendor, 0) for field, streams in streams_mapping.items() } } def get_daily_and_weekly_streams( streams_data: List[dict], current_week: DateRange, prev_week: DateRange, rolling_week: DateRange or None = None, last_date: date = None, prev_date: date = None, dsp_list: List[str] = None, market: str = None, with_vendor: bool = False, dates_and_chart_weeks: bool = True, ) -> dict: """Calculate global and market streams for the latest day, the day before, the latest week and the week before. Arguments: streams_data: Streams data. current_week: Latest week start and end dates. prev_week: Prev week start and end dates. rolling_week: Rolling week start and end dates or None. last_date: The latest date with available streams data. prev_date: Day before the latest one. dsp_list: Vendor names. market: Market code or None. with_vendor: Include vendor streams. dates_and_chart_weeks: Include dates streams. Returns: Global and market streams. """ current_week_streams = defaultdict(lambda: defaultdict(int)) prev_week_streams = defaultdict(lambda: defaultdict(int)) last_date_streams = defaultdict(dict) if dates_and_chart_weeks else None prev_date_streams = defaultdict(dict) if dates_and_chart_weeks else None rolling_week_streams = defaultdict(lambda: defaultdict(int)) if dates_and_chart_weeks else None for item in streams_data: streams = item.get("streams", 0) item_date = item["date"] vendor = item["dsp"] field_name = "global" if item["country_code"] == constants.DELPHI_GLOBAL_MARKET else "market" if item_date in current_week: current_week_streams[field_name][vendor] += streams elif item_date in prev_week: prev_week_streams[field_name][vendor] += streams if not dates_and_chart_weeks: continue if item_date == last_date.isoformat(): last_date_streams[field_name][vendor] = streams elif item_date == prev_date.isoformat(): prev_date_streams[field_name][vendor] = streams if item_date in rolling_week: rolling_week_streams[field_name][vendor] += streams streams_mapping = { "prev_week": prev_week_streams, "current_week": current_week_streams, } if dates_and_chart_weeks: streams_mapping.update( { "rolling_week": rolling_week_streams, "prev_day": prev_date_streams, "current_day": last_date_streams, } ) results = _generate_daily_weekly_streams_part(streams_mapping, "global", all_vendors=True) if market: results["market"] = market results.update(_generate_daily_weekly_streams_part(streams_mapping, "market", all_vendors=True)) if not with_vendor: return results results_vendor = {} for vendor in dsp_list: results_item = _generate_daily_weekly_streams_part(streams_mapping, "global", vendor=vendor) if market: results_item.update(_generate_daily_weekly_streams_part(streams_mapping, "market", vendor=vendor)) results_vendor[vendor] = results_item results["vendors"] = results_vendor return results def get_previous_date_range(current_range: DateRange) -> DateRange: """ Get previous date range dates based on current date values. Args: current_range (DateRange): Named tuple with start and end values of date range. Returns: DateRange: Start and end dates of previous date range. """ delta = current_range.end - current_range.start days_delta = delta.days + 1 start_extended = current_range.start - timedelta(days=days_delta) end_extended = current_range.end - timedelta(days=days_delta) return DateRange(start=start_extended, end=end_extended) def process_streams_count( data: List[dict], ) -> Tuple[Dict[str, int], Dict[str, Dict[str, int]]]: """Parse streams. Arguments: data: List of streams data. Returns: ISRC to vendor to streams count. """ streams_sum = defaultdict(int) streams_per_vendor = defaultdict(dict) for item in data: isrc = item["isrc"] streams_count = item.get("streams", 0) streams_per_vendor[item["dsp"]][isrc] = streams_count streams_sum[isrc] += streams_count return streams_sum, streams_per_vendor def collect_markets_data(data: List[dict], current_range: DateRange, prev_range: DateRange) -> List[dict]: """ Collect and aggregate track markets data for previous date range. Args: data: Streams data. current_range: Current date range values. prev_range: Previous date range values. """ markets_to_ignore = ("worldwide", "unknown", "zz") results = {} for item in data: market = item["country_code"] item_date = item["date"] if market in markets_to_ignore: continue if market not in results: current_item = { "market": market, "current_streams": 0, "prev_streams": 0, "start_date": current_range.end, "end_date": current_range.start, } results[market] = current_item else: current_item = results[market] date_obj = datetime.strptime(item_date, "%Y-%m-%d").date() if prev_range.check_in_range(date_obj): current_item["prev_streams"] += item.get("streams", 0) else: current_item["current_streams"] += item.get("streams", 0) current_item["start_date"] = min(date_obj, current_item["start_date"]) current_item["end_date"] = max(date_obj, current_item["end_date"]) results = list(results.values()) for item in results: if item["end_date"] < item["start_date"]: item["start_date"] = current_range.start item["end_date"] = current_range.end item["start_date"] = item["start_date"].isoformat() item["end_date"] = item["end_date"].isoformat() return results def process_track_in_playlist_streams_data(data: Tuple) -> List[dict]: """ Process and format raw streams data for tracks in playlist. Args: data (Tuple): Tuple of responses from Delphi with streams data. Returns: List[dict]: Track in playlist streams. """ range_fields = ["current_range", "prev_range"] tracks_dict = defaultdict(lambda: defaultdict(lambda: dict.fromkeys(range_fields, None))) for range_field, range_data in zip(range_fields, data): for market_data in range_data: isrc = market_data.get("isrc") market = market_data.get("country_code") if not isrc or not market: continue streams = market_data.get("streams", 0) tracks_dict[isrc][market][range_field] = streams items_list = [] for isrc, markets in tracks_dict.items(): markets_list = [] for market, streams in markets.items(): markets_list.append({"country_code": market, **streams}) items_list.append({"isrc": isrc, "markets": markets_list}) return items_list def process_track_playlists_streams_data(data: Tuple, vendor: str, playlists_markets: Dict[str, list]) -> List[dict]: """ Process and format Delphi streams data for track playlists in requested markets. Args: data (Tuple): Tuple of responses from Delphi with streams data. vendor (str): Vendor value (spotify/apple/amazon). playlists_markets (Dict[str, list]): Map with playlists and requested markets. Returns: List[dict]: List of playlists streams data in requested markets. """ items_dict = defaultdict( lambda: { "id": None, "country_code": None, "current_streams": None, "prev_streams": None, } ) for field, data in zip(analytics_constants.DATE_RANGES_FIELDS_NAMES, data): for item in data: playlist_id = item.get("playlist_id", "").replace(f"{vendor}_", "") country_code = item.get("country_code") if not country_code or not playlist_id or country_code not in playlists_markets.get(playlist_id, []): continue items_dict[f"{playlist_id}_{country_code}"].update( { field: item.get("streams"), "country_code": country_code, "id": playlist_id, } ) return list(items_dict.values()) def process_track_top_playlists_streams_data(date_ranges_data: List[List[dict]], vendor: str) -> dict: """ Process top tracks playlists streams data in current date range. Args: data (List[dict]): Raw Delphi response with track top playlists streams data. vendor (str): Vendor value. Returns: Tuple[dict, list]: Dict with top tracks playlists data. """ items_dict = defaultdict( lambda: { "id": None, "country_code": None, "current_streams": None, "prev_streams": None, } ) for streams_field, data in zip(analytics_constants.DATE_RANGES_FIELDS_NAMES, date_ranges_data): for item in data: vendor_playlist_id = item.get("playlist_id", "") playlist_id = vendor_playlist_id.replace(f"{vendor}_", "") # Delphi returns for Apple radios streams along with playlists one's - this condition skip them. if not playlist_id or (vendor == constants.APPLE and playlist_id.isdigit()): continue country_code, streams = item.get("country_code"), item.get("streams") key = f"{playlist_id}_{country_code}" if not country_code or streams is None: continue if streams_field == analytics_constants.PREV_STREAMS and key not in items_dict: continue items_dict[key].update( { "id": playlist_id, "country_code": country_code, streams_field: streams, } ) return items_dict def _get_youtube_item_id(ids: dict) -> str: return f"{ids.get('content_type')}_{ids.get('country_code')}_{ids.get('dsp_video_id')}" def get_short_summary_intervals( latest_date: date, ) -> Tuple[DateRange, DateRange, DateRange, DateRange]: """Get several date ranges (current day, previous day, current week, previous week). Args: latest_date: Current date. Returns: Date ranges. """ last_day = DateRange(start=latest_date, end=latest_date) last_week = DateRange(start=(latest_date - timedelta(days=6)), end=latest_date) return ( last_day, get_previous_date_range(last_day), last_week, get_previous_date_range(last_week), ) def format_short_summary(data: Tuple[List[Any], List[bool]], date_ranges: List[DateRange] or Tuple[DateRange]) -> dict: """Format short summary results data. Args: data: Short summary results. date_ranges: Date ranges. Returns: Formatted result. """ results, full_intervals = data # AP-5500: previous_week should be empty if views count for any day within this week is missing return { "current_date": results[0], "previous_date": results[1], "current_week": results[2], "previous_week": results[3] if full_intervals[3] else [], "dates": { "current_date": date_ranges[0].to_dict(isoformat=True), "previous_date": date_ranges[1].to_dict(isoformat=True), "current_week": date_ranges[2].to_dict(isoformat=True), "previous_week": date_ranges[3].to_dict(isoformat=True), }, } def calc_youtube_short_summary( data: List[dict], date_ranges: List[DateRange] or Tuple[DateRange] ) -> Tuple[List[List[dict]], List[bool]]: """Sum youtube data within date ranges, drop demographics. Args: data: Youtube data. date_ranges: List of date ranges. Returns: List of youtube data per date range and if we have data for each day for intervals or not. """ percentage_fields = ( "average_view_duration_percentage", "average_view_duration_seconds", ) data_per_range = [defaultdict(list) for _ in date_ranges] dates_per_range = [[] for _ in date_ranges] for item in data: if "dimensions" not in item: continue dimensions = item["dimensions"] item_date = dimensions.get("date") if not item_date: continue item_date = datetime.fromisoformat(item_date).date() for index, d_range in enumerate(date_ranges): if d_range.check_in_range(item_date): data_per_range[index][_get_youtube_item_id(dimensions)].append(item) dates_per_range[index].append(item_date.toordinal()) results = [[] for _ in date_ranges] for index, range_items in enumerate(data_per_range): for items_group in range_items.values(): result_item = {} results[index].append(result_item) result_percentage = {f: 0 for f in percentage_fields} for item in items_group: utils.deep_sum(result_item, item) metrics = item.get("metrics", {}) for field in percentage_fields: result_percentage[field] += metrics.get(field, 0) * metrics.get("views", 0) metrics = result_item.get("metrics", {}) views = metrics["views"] for field in percentage_fields: if field in metrics: metrics[field] = round(result_percentage[field] / views, 6) if views > 0 else 0 del result_item["dimensions"]["date"] return results, [len(set(dates)) == date_ranges[index].days for index, dates in enumerate(dates_per_range)] def calc_tiktok_short_summary( data: Dict[str, Dict[str, Dict[str, List[int]]]] ) -> Tuple[List[Dict[str, Dict[str, Dict[str, int]]]], List[bool]]: """Sum tiktok data within date ranges. Args: data: Tiktok data. Returns: Tiktok data per date range and flags if we have data for each day in intervals or not. """ # We have data here for the last two weeks, Delphi API returns this as dict (metric name to list). # The lists include metrics for each day from min date to max. So two weeks => several lists of 14 items within. # Here we define several date ranges as indexes in the lists: last day, previous day, last week and previous week. date_ranges = ((-1, None), (-2, -1), (-7, None), (0, -7)) dates_in_range = [abs(r[1] or 0 - r[0] or 0) for r in date_ranges] data_per_range = [defaultdict(lambda: defaultdict(dict)) for _ in date_ranges] all_dates = [True] * len(date_ranges) for content_type_name, content_type_values in data.items(): for country_code_name, country_code_values in content_type_values.items(): for field, values in country_code_values.items(): for index, date_range in enumerate(date_ranges): range_values = values[date_range[0] : date_range[1]] if len(range_values) < dates_in_range[index] or None in range_values: all_dates[index] = False data_per_range[index][content_type_name][country_code_name][field] = sum( i or 0 for i in range_values ) return data_per_range, all_dates def process_playlists_streams_for_graph(data: List[dict], date_range: DateRange, vendor: str) -> List[dict]: """Process and format Delphi playlists streams data. Args: data (List[dict]): Delphi response with playlists streams in specified markets. date_range (DateRange): Date range object. vendor (str): Vendor value. Returns: List[dict]: List of playlists graph coordinates data. """ playlist_market_map = defaultdict(lambda: dict()) for item in data: playlist_id, country_code = item.get("playlist_id", "").replace(f"{vendor}_", ""), item.get("country_code") if not playlist_id or not country_code: continue item_date, streams = ( datetime.strptime(item["date"], "%Y-%m-%d"), item["streams"], ) playlist_market_map[f"{playlist_id}__{country_code}"][utils.date_to_timestamp(item_date)] = streams result = [] for pl_id_market, streams_data in playlist_market_map.items(): playlist_id, market = pl_id_market.split("__") result.append( { "country_code": market, "playlist_id": playlist_id, "coordinates": _generate_coordinates(streams_data, date_range.start, date_range.end, default_none=True), } ) return result def calc_trend(current_value: int or None, previous_value: int or None) -> int or None: """Calculate trend. Args: current_value: The latest value. previous_value: The previous value. Returns: Trend. """ return ( round((current_value - previous_value) / previous_value * 100) if current_value is not None and previous_value else None ) def calc_period_trend( current_value: int or None, previous_value: int or None, current_days: int, previous_days: int, period: int = 7, ) -> int or None: """Calculate week trend. Args: current_value: The latest week value. previous_value: The previous week value. current_days: The latest week not None days count. previous_days: The previous week not None days count. period: Returns: Trend. """ return None if current_days < period or previous_days < period else calc_trend(current_value, previous_value) async def get_metrics_trends_data( api: DelphiClient, isrc_list: List[str], market: str, metric_list: List[str], start_date: date, end_date: date, ) -> Dict[str, List[int]]: """Get metrics trends data. Args: api: Delphi API client. isrc_list: ISRC list. market: Market code. metric_list: Metric list. start_date: Two week range start. end_date: Two week range end. Returns: Metrics data. """ dsp_list = {analytics_constants.MetricsTrendsValuesMetric.DSP_MAPPING[i] for i in metric_list} tasks = [] if constants.SPOTIFY in dsp_list: tasks.append( api.get_streams( start_date=start_date, end_date=end_date, isrc_list=isrc_list, dsp_list=[constants.SPOTIFY], country_code_list=[market], group_by=[StreamsGroupBy.DATE], sort_by=delphi_constants.SortBy.DATE, sort_order=SortOrder.ASC, ) ) if constants.TIKTOK in dsp_list: tasks.append( api.get_tiktok_tracks_analytics( start_date=start_date, end_date=end_date, isrc_list=isrc_list, country_code_list=[market], breakdowns_list=[TiktokTrackAnalyticsBreakdown.DAILY], metrics_list=[ i for i in metric_list if analytics_constants.MetricsTrendsValuesMetric.DSP_MAPPING[i] == constants.TIKTOK ], use_args_dates=True, ) ) responses = await asyncio.gather(*tasks) metrics_data = {} if constants.SPOTIFY in dsp_list: views_data = [] if responses[0]: date_streams = defaultdict(int) for item in responses[0]: date_streams[item["date"]] += item["streams"] for day in range(14): current_date = (start_date + timedelta(days=day)).isoformat() views_data.append(date_streams[current_date] if current_date in date_streams else None) metrics_data[analytics_constants.MetricsTrendsValuesMetric.STREAMS] = views_data if constants.TIKTOK in dsp_list: metrics_data.update(responses[-1]["breakdowns"][TiktokTrackAnalyticsBreakdown.DAILY]) return metrics_data def calc_metrics_trends( metrics_data: Dict[str, List[int]], metric_list: List[str], start_date: date, end_date: date, period: int = 7, ) -> List[dict]: """Calculate metrics trends. Args: metrics_data: Metric to values list mapping. metric_list: List of metrics. start_date: The latest week start date. end_date: The latest date (week end date). period: value to compare with data existence for period Returns: Metrics trends. """ result = [] period_in_days = (end_date - start_date).days + 1 double_period = period_in_days * 2 for metric in metric_list: values = metrics_data[metric] # metric values for 2 periods count = len(values) latest_value = ( values[double_period - 1] if count == double_period else None ) # the latest day value is the last in the list period_ago_value = values[period_in_days - 1] if count >= period_in_days else None # value for a period before the latest day one_day_ago_value = values[-2] if count >= period_in_days else None period_values = values[period_in_days:] # the last period values period_days = sum(0 if i is None else 1 for i in period_values) period_sum = sum(i or 0 for i in period_values) if period_values and period_days > 0 else None period_avg = round(period_sum / period_days) if period_values and period_days else None previous_period_values = values[:period_in_days] previous_period_days = sum(0 if i is None else 1 for i in previous_period_values) previous_period_sum = ( sum(i or 0 for i in previous_period_values) if previous_period_values and previous_period_days > 0 else None ) previous_period_avg = ( round(previous_period_sum / previous_period_days) if previous_period_values and previous_period_days else None ) result.append( { "dsp": analytics_constants.MetricsTrendsValuesMetric.DSP_MAPPING[metric], "metric": metric, "latest": { "date": end_date, "value": latest_value, # (({qty for the latest day} - {qty for period days ago}) / {qty for period days ago}) * 100% "trend": calc_trend(latest_value, period_ago_value), # (({qty for the latest day} - {qty for the day 1 days ago}) / {qty for the day 1 days ago}) * 100% "trend_one_day_ago": calc_trend(latest_value, one_day_ago_value), }, # TODO: need to change week to period during next refactoring "week": { "start_date": start_date, "end_date": end_date, "value": period_sum, # (({qty for the latest prd} - {qty for the previous prd}) / {qty for the previous prd}) * 100% "trend": calc_period_trend( period_sum, previous_period_sum, period_days, previous_period_days, period ), "average": { "value": period_avg, # (({Avg for the latest prd} - {Avg for the previous prd}) / {Avg for the previous prd}) * 100% "trend": calc_period_trend( period_avg, previous_period_avg, period_days, previous_period_days, ), }, "data": period_values, }, } ) return result async def get_metrics_markets_week_values_data( api: DelphiClient, isrc_list: List[str], market: str or None, metric_list: List[str], start_date: date, end_date: date, ) -> Dict[str, Dict[str, int]]: """Get metrics markets week values data. Args: api: Delphi API client. isrc_list: ISRC list. market: Market code or None for top. metric_list: Metric list. start_date: Week start date. end_date: Week end date. Returns: Metrics markets week values. """ dsp_list = {analytics_constants.MetricsTrendsValuesMetric.DSP_MAPPING[i] for i in metric_list} market_list = [market, constants.DELPHI_GLOBAL_MARKET] if market else None tasks = [] if constants.SPOTIFY in dsp_list: tasks.append( api.get_streams( start_date=start_date, end_date=end_date, isrc_list=isrc_list, dsp_list=[constants.SPOTIFY], country_code_list=market_list, ) ) if constants.TIKTOK in dsp_list: tasks.append( api.get_tiktok_tracks_analytics( start_date=start_date, end_date=end_date, isrc_list=isrc_list, country_code_list=market_list, breakdowns_list=[TiktokTrackAnalyticsBreakdown.COUNTRY_TOTALS], metrics_list=[ i for i in metric_list if analytics_constants.MetricsTrendsValuesMetric.DSP_MAPPING[i] == constants.TIKTOK ], ) ) responses = await asyncio.gather(*tasks) metrics_data = defaultdict(lambda: defaultdict(int)) if constants.SPOTIFY in dsp_list: for item in responses[0]: metrics_data[analytics_constants.MetricsTrendsValuesMetric.STREAMS][item["country_code"]] += item["streams"] if constants.TIKTOK in dsp_list: breakdown_data = responses[-1]["breakdowns"][TiktokTrackAnalyticsBreakdown.COUNTRY_TOTALS] for country_code, country_values in breakdown_data.items(): for metric_name, metric_value in country_values.items(): metrics_data[metric_name][country_code] = metric_value return metrics_data def calc_metrics_markets_week_values( metrics_data: Dict[str, Dict[str, int]], market: str or None, start_date: date, end_date: date, ) -> List[dict]: """Calculate metrics markets week values. Args: metrics_data: Metric to values list mapping. market: Selected market or None for top. start_date: Week start date. end_date: Week end date. Returns: Metrics markets week values. """ result = [] for metric, markets_values in metrics_data.items(): if market: metric_market = market else: metric_market = max( markets_values, key=lambda key: -1 if key == constants.DELPHI_GLOBAL_MARKET else markets_values[key], ) worldwide_value = markets_values[constants.DELPHI_GLOBAL_MARKET] market_value = markets_values[metric_market] result.append( { "dsp": analytics_constants.MetricsTrendsValuesMetric.DSP_MAPPING[metric], "metric": metric, "week": { "start_date": start_date, "end_date": end_date, "value": { constants.DELPHI_GLOBAL_MARKET: worldwide_value, "market": market_value, }, "market": metric_market, "percentage": round(market_value / worldwide_value * 100) if worldwide_value else None, }, } ) return result def get_shazam_chart_id(market: str, city: str or None) -> str: """Get shazam chart ID. Args: market: Market code. city: City ID or None. Returns: Chart ID. """ return f"shazam_{market}{f'_{city}' if city else ''}" def calc_tiktok_insight(metrics_data: Dict[str, List[int]], start_date: date) -> Dict: creations = metrics_data["creations"] video_views = metrics_data["video_views"] peak_creation = max([val for val in creations if val is not None]) peak_creation_indices = [i for i, x in enumerate(creations) if x == peak_creation] peak_creation_views = {view: i for i, view in enumerate(video_views) if i in peak_creation_indices} peak_creation_view = max([key for key in peak_creation_views.keys() if key is not None]) final_max_index = video_views.index(peak_creation_view) peak_date = start_date + timedelta(days=final_max_index) return { "creations": peak_creation, "views": peak_creation_view, "peak_date": peak_date, } def get_tiktok_total(total_data: Dict) -> Dict: totals = get_combined_data([{"key": "key", **val} for val in total_data.values()], group_keys=["key"]) result = totals[0] result.pop("key") return result def build_tiktok_graph(tiktok_data: Dict, start_date: date) -> Dict: metrics_data = tiktok_data["breakdowns"][TiktokTrackAnalyticsBreakdown.DAILY] totals = tiktok_data["breakdowns"][TiktokTrackAnalyticsBreakdown.TOTALS] totals = get_tiktok_total(totals) creations = metrics_data["creations"] video_views = metrics_data["video_views"] creations_graph = [] video_views_graph = [] for i, value in enumerate(creations): point_date = start_date + timedelta(days=i) x = int(datetime.combine(point_date, datetime.min.time()).replace(tzinfo=timezone.utc).timestamp()) creations_graph.append({"x": x, "y": value}) video_views_graph.append({"x": x, "y": video_views[i]}) return { "creations": {"graph": creations_graph, "total": totals["creations"]}, "video_views": {"graph": video_views_graph, "total": totals["video_views"]}, } def calc_tiktok_top(country_codes_data: Dict, limit: int) -> Dict: totals = country_codes_data.pop("worldwide") country_codes_data = [{"market": code, **val} for code, val in country_codes_data.items()] result = {} for metric in TIKTOK_METRICS_LIST: if not totals[metric]: result[metric] = None break country_codes_data = multikeysort(country_codes_data, [f"-{metric}", "market"]) creations_top = country_codes_data[:limit] result_list = [] for row in creations_top: # AG-10385 filter zero value markets if row[metric] > 0: result_list.append( { "market": row["market"], "cnt": row[metric], "percentage": round(row[metric] / totals[metric] * 100), } ) result[metric] = result_list return result def has_tiktok_none_data(tiktok_data: Dict): metrics_data = tiktok_data["breakdowns"][TiktokTrackAnalyticsBreakdown.DAILY] creations = metrics_data["creations"] video_views = metrics_data["video_views"] is_none = creations or video_views return is_none def get_amazon_playlist_worldwide_streams_graph( streams_data: List[Dict[str, Any]], interested_playlists: List[str], group_playlist_id: str, group_market: str ) -> List[Dict[str, Any]]: """Aggregate streams data by date for list of playlists Args: streams_data: Streams data interested_playlists: Playlists we are interested in group_playlist_id: response name for grouped objects group_market: grouped objects country_code name Returns: Aggregated streams result """ result = defaultdict(dict) playlists_to_check = [f"{pl_id[:pl_id.index(':')]}_{pl_id[pl_id.index('_')+1:]}" for pl_id in interested_playlists] for data in streams_data: if f"{data['playlist_id'].replace(f'{DSP.AMAZON.value}_','')}_{data['country_code']}" in playlists_to_check: if data["date"] not in result.keys(): data["country_code"] = group_market data["playlist_id"] = group_playlist_id result[data["date"]] = data continue result[data["date"]]["streams"] += data["streams"] return result.values()