import glom from abc import abstractmethod from collections import defaultdict from datetime import date from typing import Any, Callable, Dict, Iterable, List, Optional, Tuple from server.legacy.core.constants import DELPHI_GLOBAL_MARKET, GLOBAL_MARKET_CODE from server.legacy.core.utils import convert_market from server.legacy.core.vendor_client import BaseVendor as Vendor class BaseVendor(Vendor): """Base vendor class for Consumer Analytics client.""" vendor: str @classmethod def get_streams_info(cls, item): return item.get(f"{cls.vendor}_streams_info", {}) @staticmethod def to_camel_case(field: str, prefix: str = ""): parts = field.split("_") if prefix: return prefix + "".join(x.capitalize() for x in parts) else: return "".join([parts[0]] + [x.capitalize() for x in parts[1:]]) def parse_market_item(self, market_item, date_item_keys: Iterable = None, **kwargs): return [ self.parse_date_item(_date, date_item, date_item_keys, **kwargs) for _date, date_item in sorted(market_item.items()) ] def parse_isrc_item_v0( self, isrc: str, isrc_item: Dict[date, dict], data_item_keys: Iterable = None, **kwargs ) -> Optional[dict]: """Parse streams and source types data item for particular isrc from compact format to v0 format view. Args: isrc: ISRC code. isrc_item: data to parse. data_item_keys: tuple of flat and packed keys in item. Returns: Track per country data. """ return { "isrc": isrc, "data": [ {"countryCode": market, "data": self.parse_market_item(market_item, data_item_keys, **kwargs)} for market, market_item in sorted(isrc_item.items()) if market_item ], } def parse_isrc_item_v1( self, isrc: str, isrc_item: Dict[date, dict], market: str, data_item_keys: Tuple[tuple, dict] = None ) -> Optional[dict]: """Parse streams and source types data item for particular isrc from compact format dict to v1 format view with filtering by passed market. Args: isrc: ISRC code. isrc_item: data to parse. market: selected market to get data for, required. data_item_keys: tuple of flat and packed keys in item. Returns: list of { "isrc": particular item isrc, "data": list of parsed date item maps, each based on particular vendor format. } """ if market is None: raise ValueError("Version #1 is available for selected market only.") market_item = isrc_item.get(market) if market_item: return {"isrc": isrc, "data": self.parse_market_item(market_item, data_item_keys)} @classmethod def map_data_by_isrc( cls, streams_data: List[dict], data_item_keys: Tuple[tuple, dict], market: str = None, exclude_global: bool = True, ) -> Dict[str, Dict[str, Dict[date, dict]]]: """Group items by ISRC and do clean up of excess data. Args: streams_data: several isrc data to parse. data_item_keys: tuple of flat and packed keys in item. market: selected market to get data only for, use None for all markets. exclude_global: Flag to exclude global (_gl) data from output. Returns: mapped by ISRC and market streaming data. """ flat_keys, packed_keys = data_item_keys streams_data_map = defaultdict(lambda: defaultdict(dict)) for item in streams_data: item_market = item["country_code"] if (market and item_market != market) or (exclude_global and item_market == DELPHI_GLOBAL_MARKET): continue if item_market == DELPHI_GLOBAL_MARKET: item_market = GLOBAL_MARKET_CODE streams_info = cls.get_streams_info(item) isrc = item["isrc"] item_date = item.get("date") result_item = {k: item.get(k) or streams_info.get(k, 0) for k in flat_keys} for k in packed_keys.keys(): result_item[k] = streams_info.get(k, {}) streams_data_map[isrc][item_market][item_date] = result_item return dict(streams_data_map) def parse_isrc_map( self, isrc_list: List[str], streams_data: List[dict], market: str = None, streams_only: bool = False, exclude_global: bool = True, combine_isrc: bool = False, item_parser: Callable = parse_isrc_item_v1, **kwargs, ): """Parse streams and source types data map for several isrc, combine data for all isrc if needed. Args: isrc_list: list of all requested ISRC. streams_data: several isrc data to parse. market: selected market to get data only for, use None for all markets. streams_only: flag, if True return streams data only without sources types, etc. combine_isrc: flag, if True return combined (summed up) data for all isrc in given list. item_parser: function to parse one isrc item, for now parse_isrc_item_v1, parse_isrc_item_v0 are supported. exclude_global: Flag to exclude global (_gl) data from output. Returns: list of { "isrc": particular item isrc, "data": list of data dictionaries in format defined by item_parser. }. """ data_item_keys = self.get_data_item_keys(streams_only) streams_data = self.map_data_by_isrc(streams_data, data_item_keys, market, exclude_global) if combine_isrc: streams_data = { ",".join(isrc_list): self.get_combined_by_isrc_map(list(streams_data.values()), data_item_keys) } result = [] for isrc, isrc_item in sorted(streams_data.items()): parsed_item = item_parser(isrc, isrc_item, market=market, data_item_keys=data_item_keys, **kwargs) if parsed_item: result.append(parsed_item) return result @classmethod def get_combined_by_isrc_map( cls, streams_data: List[Dict[str, Dict[date, dict]]], data_item_keys: Tuple[tuple, dict] ) -> Dict[str, Dict[date, dict]]: """Parse streams and source types data map for several isrc, combine isrc data by markets iand dates. Args: streams_data: several isrc data to parse. data_item_keys: iterable of item keys. Returns: combined data. """ if not streams_data: return {} flat_keys, packed_keys = data_item_keys result = streams_data[0] for market_item in streams_data[1:]: for market, date_items in market_item.items(): if market not in result: result[market] = date_items continue result_market_item = result[market] for item_date, item in date_items.items(): if item_date not in result_market_item: result_market_item[item_date] = item continue result_date_item = result_market_item.get(item_date, {}) for k in flat_keys: result_date_item[k] = result_date_item.get(k, 0) + item.get(k, 0) for node, fields in packed_keys.items(): if node not in result_date_item: result_date_item[node] = {} result_date_item_node = result_date_item[node] data_item_node = item.get(node, {}) for field in fields: result_date_item_node[field] = result_date_item_node.get(field, 0) + data_item_node.get( field, 0 ) return result @staticmethod @abstractmethod def parse_date_item(_date: date, date_item: Dict[str, Any], keys: Iterable, **kwargs) -> Dict[str, Any]: """Parse streams and source types data map for particular date from compact to full format. Args: _date: date key. date_item: data map to parse. keys: iterable of keys to get data from passed data map by. Returns: data map for one date based on particular vendor format. """ pass @staticmethod @abstractmethod def get_data_item_keys(streams_only: bool = False) -> Tuple[tuple, dict]: """Return tuple of flat keys, packed keys iterable for date streams item.""" pass @classmethod def convert_to_track_compat_v2(cls, streams_data: List[dict]) -> dict: """Convert Delphi streams response to CA track compat v2. Args: streams_data: Delphi streams data. Returns: CA like track compat v2 streams. """ results = defaultdict(lambda: defaultdict(dict)) for item in streams_data: streams_info = cls.get_streams_info(item) if not streams_info: continue results_item = cls.convert_item_to_track_compat_v2(item.get("streams", 0), streams_info) market = convert_market(item["country_code"], GLOBAL_MARKET_CODE) results[item["isrc"]][market][item["date"]] = results_item return {k: dict(sorted(v.items())) for k, v in results.items()} @classmethod def convert_to_track_compat_v1(cls, streams_data: List[dict]) -> List[dict]: """Convert Delphi streams response to CA track compat v1. Args: streams_data: Delphi streams data. Returns: CA like track compat v1 streams. """ result = defaultdict(dict) for item in streams_data: streams_info = cls.get_streams_info(item) market = item["country_code"] item_date = item["date"] if market == DELPHI_GLOBAL_MARKET: continue result[market][item_date] = cls.convert_item_to_track_compat_v1( item.get("streams", 0), streams_info, item_date, result[market].get(item_date) ) return [{"countryCode": market, "data": list(dates.values())} for market, dates in sorted(result.items())] @classmethod @abstractmethod def convert_item_to_track_compat_v2(cls, streams_count: int, streams_info: dict) -> dict: pass @classmethod @abstractmethod def convert_item_to_track_compat_v1( cls, streams_count: int, streams_info: dict, item_date: date, item: dict ) -> dict: pass @classmethod def parse_tracks_playlists_streams( cls, data: List[dict], market: str, include_shuffle: bool = False, per_day: bool = False, include_summary: bool = False, all_markets: bool = False, add_isrc: bool = False, ) -> dict: """Parse tracks playlists streams per day. Args: data: Delphi streams data. market: Selected market to consider as local. include_shuffle: Include shuffle streams. per_day: Stats per day or summary. include_summary: Include summary block for per day streams. all_markets: Flag to parse additional markets. add_isrc: Flag fo adding 'isrc' key. Returns: CA like result. """ stats = defaultdict(dict) summary = defaultdict(int) for data_item in data: item_market = data_item["country_code"] is_global = item_market == DELPHI_GLOBAL_MARKET market_type = is_global and "g" or (item_market == market) and "l" or None streams_count = data_item.get("streams", 0) stats_item = ( stats[data_item["date"]] if per_day else stats[data_item["playlist_id"].replace(f"{cls.vendor}_", "")] ) if all_markets and not is_global: days_spec = f"markets.{item_market}.d" glom.assign(stats_item, days_spec, glom.glom(stats_item, days_spec, default=0) + 1, missing=dict) streams_spec = f"markets.{item_market}.st" glom.assign( stats_item, streams_spec, glom.glom(stats_item, streams_spec, default=0) + streams_count, missing=dict, ) if market_type is None: continue stats_item[f"{market_type}_st"] = stats_item.get(f"{market_type}_st", 0) + streams_count summary[f"{market_type}_st"] += streams_count if add_isrc: stats_item.setdefault("isrc", data_item.get("isrc")) if not per_day: stats_item[f"{market_type}_d"] = stats_item.get(f"{market_type}_d", 0) + 1 if include_shuffle: shuffle_count = cls.get_streams_info(data_item).get("shuffle_play", 0) stats_item[f"{market_type}_sh"] = stats_item.get(f"{market_type}_sh", 0) + shuffle_count summary[f"{market_type}_sh"] += shuffle_count stats = dict(sorted(stats.items())) market_field = "market" if include_summary else "local" if per_day: results_items = cls.generate_tracks_playlists_per_day_result( stats, include_shuffle, market_field=market_field ) else: results_items = cls.generate_tracks_playlists_summary_result( stats, include_shuffle, all_markets=all_markets, add_isrc=add_isrc ) if per_day and include_summary: results = cls.generate_tracks_playlists_per_day_item(None, summary, include_shuffle, market_field) results["values"] = results_items return results return {"items": results_items} @staticmethod @abstractmethod def generate_tracks_playlists_summary_result( data: Dict[str, Dict[str, int]], include_shuffle: bool, all_markets: bool = False, add_isrc: bool = False ) -> list: """Generate CA like tracks playlists streams summary result. Args: data: Parsed stats. include_shuffle: Include shuffle streams. all_markets: Parse additional markets add_isrc: Add 'isrc' attribute. Returns: CA like tracks playlists streams summary result. """ pass @staticmethod @abstractmethod def generate_tracks_playlists_per_day_item( streams_date: date or None, data: Dict[str, int], include_shuffle: bool, market_field: str = "local" ) -> dict: """Generate CA like tracks playlists streams per day result. Args: streams_date: Streams date. data: Parsed stats. include_shuffle: Include shuffle streams. market_field: Market fields prefix. Returns: CA like tracks playlists streams per day result item. """ pass @classmethod def generate_tracks_playlists_per_day_result( cls, data: Dict[date, Dict[str, int]], include_shuffle: bool, market_field: str = "local" ) -> list: """Generate CA like tracks playlists streams per day result. Args: data: Parsed stats. include_shuffle: Include shuffle streams. market_field: Market fields prefix. Returns: CA like tracks playlists streams per day result. """ return [ cls.generate_tracks_playlists_per_day_item(streams_date, streams_item, include_shuffle, market_field) for streams_date, streams_item in data.items() ] @classmethod def parse_streams_from_release_to_date(cls, streams_data: List[dict], date_data: List[dict]) -> dict: """Parse streams and first date data and return in CA view. Args: streams_data: Track streaming data. date_data: Track first streaming data. Returns: CA like streams from release to date response. """ return { "firstStreamDate": date_data[0].get("date") if date_data else None, "items": [ cls.parse_streams_from_release_to_date_item( convert_market(item["country_code"], GLOBAL_MARKET_CODE), item.get("streams", 0), cls.get_streams_info(item), ) for item in streams_data ], } @classmethod @abstractmethod def parse_streams_from_release_to_date_item(cls, market: str, streams_count: int, streams_info: dict) -> dict: pass