import math from collections import defaultdict from datetime import date import json import inspect import os from typing import Iterable, List from unittest.mock import patch import pytest from server.constants.core import DELPHI_GLOBAL_MARKET 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 import dsp from server.constants.delphi.videos.group_by import VideosGroupBy from server.constants.delphi.videos.only import VideosOnly from server.constants.sorting import SortOrder from server.client.clients import delphi_client from server.legacy.delphi.constants import START_DATE_MIN from server.utils.delphi.converters import str_to_date from server.utils.delphi.misc import nullable_max, nullable_min streams_args_spec = inspect.getfullargspec(delphi_client.DelphiClient.get_streams).args[1:] charts_args_spec = inspect.getfullargspec(delphi_client.DelphiClient.get_charts).args[1:] first_dates_args_spec = inspect.getfullargspec(delphi_client.DelphiClient.get_first_stream_dates).args[1:] video_analytics_args_spec = inspect.getfullargspec(delphi_client.DelphiClient.get_video_analytics).args[1:] tiktok_tracks_analytics_args_spec = inspect.getfullargspec( delphi_client.DelphiClient._get_tiktok_tracks_analytics ).args[1:] regions_args_spec = inspect.getfullargspec(delphi_client.DelphiClient.get_regions).args[1:] get_args_spec = inspect.getfullargspec(delphi_client.DelphiClient._get).args[1:] def read_json_from_file(filename: str) -> list or dict: """Read json from file.""" with open(f"{os.path.dirname(__file__)}/{filename}.json", "r") as f: data = f.read() return json.loads(data) def _get_date_from_kwargs(field: str, kwargs: dict) -> date: result = kwargs[field] if isinstance(result, str): result = date.fromisoformat(result) return result def _skip_param(item: dict, kwargs: dict, field: str, arg_name: str, container_field: str = None): return ( kwargs.get(arg_name) and (item[container_field][field] if container_field else item[field]) not in kwargs[arg_name] ) streams_min_date = date(2020, 10, 1) streams_max_date = date(2020, 10, 15) video_analytics_min_date = date(2021, 1, 15) video_analytics_max_date = date(2021, 1, 31) async def get_streams(*args, **kwargs) -> list: """Emulate get streams Delphi endpoint.""" kwargs = kwargs or {} kwargs.update(zip(streams_args_spec, args)) start_date = _get_date_from_kwargs("start_date", kwargs) end_date = _get_date_from_kwargs("end_date", kwargs) start_date_str = start_date.isoformat() end_date_str = end_date.isoformat() if start_date > streams_max_date or end_date < streams_min_date: return [] group_by = kwargs.get("group_by") if group_by and StreamsGroupBy.DATE in group_by: postfix = "by_date_by_sub_dsp" if StreamsGroupBy.SUB_DSP in group_by else "by_date" else: postfix = "summary" group_by_date = kwargs.get("group_by") == [StreamsGroupBy.DATE] subset_playlists = kwargs.get("subset") == StreamsSubset.PLAYLISTS is_from_first_streams_date = start_date == START_DATE_MIN isrc_list = kwargs.get("isrc_list") if not subset_playlists: filename = f"tracks_{postfix}" elif subset_playlists and isrc_list: filename = f"track_playlists_{postfix}" else: filename = f"playlists_{postfix}" data = read_json_from_file(filename) results = [] for item in data: # filter data if ( _skip_param(item, kwargs, "isrc", "isrc_list") or _skip_param(item, kwargs, "dsp", "dsp_list") or _skip_param(item, kwargs, "country_code", "country_code_list") ): continue if group_by_date and (item["date"] < start_date_str or item["date"] > end_date_str): continue if ( subset_playlists and kwargs.get("playlist_id_list") and item["playlist_id"] not in kwargs["playlist_id_list"] ): continue results.append(item) # for from first streams date to selected increase streams count x10 if is_from_first_streams_date: item["streams"] *= 10 if not group_by_date and (start_date > streams_min_date or end_date < streams_max_date): item["streams"] = int( item["streams"] * (min(end_date, streams_max_date) - max(start_date, streams_min_date)).days / 8 ) # handle include include = kwargs.get("include") or [] if StreamsInclude.ALL in include: continue vendor = item["dsp"] if vendor.startswith(dsp.AMAZON): vendor = dsp.AMAZON if StreamsInclude.DEMOGRAPHICS not in include and vendor != dsp.AMAZON: if "genders" in item: del item["genders"] if f"{vendor}_age_bands" in item: del item[f"{vendor}_age_bands"] if not include: del item[f"{vendor}_streams_info"] continue streams_info = item[f"{vendor}_streams_info"] if StreamsInclude.ENGAGEMENT not in include: del streams_info["listener_engagement" if vendor == dsp.APPLE else "engagement"] if vendor == dsp.AMAZON: continue if StreamsInclude.SKIPS not in include: del streams_info["skips"] if StreamsInclude.SAVES not in include and "saves" in streams_info: del streams_info["saves"] if StreamsInclude.SOURCES not in include: del streams_info["source"] sort_by, sort_order = kwargs.get("sort_by"), kwargs.get("sort_order", SortOrder.ASC) if sort_by: results = list(sorted(results, key=lambda i: i[sort_by], reverse=(sort_by == SortOrder.DESC))) return results isrc_to_video_id = {"GBARL1401524": "OPf0YbXqDm0", "BRSME1600749": "YJ6F8qlINbc"} VIDEOS_FIELDS_MAPPING = { VideosOnly.DEMOGRAPHICS: ("views", "gender_views_percentages", "age_band_views_percentages"), VideosOnly.VIEWS: ("views",), VideosOnly.WATCH_TIMES: ("views", "average_view_duration_percentage", "average_view_duration_seconds"), VideosOnly.DEVICES: ("views", "device_type", "operating_system"), VideosOnly.ENGAGEMENT: ("views", "comments", "shares", "likes", "dislikes"), VideosOnly.TRAFFIC_SOURCES: ("views", "traffic_source_types"), } async def get_video_analytics(*args, **kwargs) -> list: """Emulate get video analytics Delphi endpoint.""" kwargs = kwargs or {} kwargs.update(zip(video_analytics_args_spec, args)) start_date = _get_date_from_kwargs("start_date", kwargs) end_date = _get_date_from_kwargs("end_date", kwargs) start_date_str = start_date.isoformat() end_date_str = end_date.isoformat() isrc = kwargs.get("isrc") if isrc and not isinstance(isrc, (list, tuple, set)): isrc = [isrc] video_id = [isrc_to_video_id.get(i) for i in isrc] if isrc else kwargs.get("video_id") if video_id: if not isinstance(video_id, (list, tuple, set)): video_id = [video_id] video_id = [i.replace("youtube_", "") for i in video_id] country_code_list = kwargs.get("country_code_list") only = kwargs.get("only") if start_date > video_analytics_max_date or end_date < video_analytics_min_date: return [] group_by = kwargs.get("group_by", []) if group_by and VideosGroupBy.DATE in group_by and VideosGroupBy.VIDEO_ID in group_by: postfix = "by_date_video" else: postfix = "summary" group_by_date = StreamsGroupBy.DATE in group_by filename = f"video_analytics_{postfix}" data = read_json_from_file(filename) results = [] for item in data: # filter data if _skip_param(item, kwargs, "content_type", "content_type", "dimensions") or ( video_id and not item["dimensions"]["dsp_video_id"] in video_id ): continue country_code = item["dimensions"]["country_code"] if country_code_list and country_code not in country_code_list: continue if group_by_date and (item["dimensions"]["date"] < start_date_str or item["dimensions"]["date"] > end_date_str): continue if only: metrics = item["metrics"] for only_key, only_fields in VIDEOS_FIELDS_MAPPING.items(): if only_key not in only: for field in only_fields: if field != "views": metrics.pop(field, None) if not metrics: continue results.append(item) return results first_dates_data = [ {"country_code": "worldwide", "date": "2020-02-07", "dsp": "amazon", "isrc": "USSM19902989", "streams": 1549}, {"country_code": "worldwide", "date": "2019-11-05", "dsp": "apple", "isrc": "USSM19902989", "streams": 222708}, {"country_code": "worldwide", "date": "2019-11-06", "dsp": "spotify", "isrc": "USSM19902989", "streams": 771218}, {"country_code": "worldwide", "date": "2020-02-08", "dsp": "amazon", "isrc": "USSM19902991", "streams": 6813}, {"country_code": "worldwide", "date": "2019-11-03", "dsp": "apple", "isrc": "USSM19902991", "streams": 144907}, {"country_code": "worldwide", "date": "2019-11-04", "dsp": "spotify", "isrc": "USSM19902991", "streams": 446618}, {"country_code": "worldwide", "date": "2020-02-09", "dsp": "amazon", "isrc": "USSM19902990", "streams": 7028}, {"country_code": "worldwide", "date": "2019-11-01", "dsp": "apple", "isrc": "USSM19902990", "streams": 79689}, {"country_code": "worldwide", "date": "2019-11-02", "dsp": "spotify", "isrc": "USSM19902990", "streams": 300590}, ] async def get_first_streams_dates(*args, **kwargs): """Emulate get first streams dates Delphi endpoint.""" kwargs = kwargs or {} kwargs.update(zip(first_dates_args_spec, args)) data = list(first_dates_data) return [ item for item in data if not _skip_param(item, kwargs, "isrc", "isrc_list") and not _skip_param(item, kwargs, "dsp", "dsp_list") ] latest_dates_data = [ {"dsp": "spotify", "updated_date": "2020-10-15"}, {"dsp": "apple", "updated_date": "2020-10-16"}, {"dsp": "amazon", "updated_date": "2020-10-14"}, {"dsp": "youtube", "updated_date": "2021-01-31"}, {"dsp": "tiktok", "updated_date": "2020-10-17"}, ] def sum_data(item1: dict, item2: dict, wrap_none: bool = True): for key, value in item2.items(): if isinstance(value, dict): item1[key] = sum_data(item1.get(key, {}), value, wrap_none=wrap_none) elif value is None and wrap_none or isinstance(value, int): item1[key] = (item1.get(key) or 0) + (value or 0) elif value is None and not wrap_none: item1.setdefault(key, None) else: item1[key] = value return item1 def combine_vendor_data(data: List[dict], isrc_list: List[str], wrap_none: bool = True) -> List[dict]: results = defaultdict(lambda: defaultdict(lambda: defaultdict(dict))) combined_isrc = ",".join(isrc_list) for item in data: combined_item = results[item["dsp"]][item["country_code"]][item.get("date")] sum_data(combined_item, item, wrap_none=wrap_none) combined_item["isrc"] = combined_isrc return [ date_item for dsp_item in results.values() for market_item in dsp_item.values() for date_item in market_item.values() ] def combine_data_by_tiers(data: List[dict], wrap_none: bool = True, vendor: str = dsp.AMAZON) -> List[dict]: results = defaultdict(lambda: defaultdict(lambda: defaultdict(dict))) for item in data: combined_item = results[item["isrc"]][item["country_code"]][item["date"]] sub_dsp, sub_dsp_streams = item["dsp"], item["streams"] if vendor: sub_dsp = sub_dsp[len(vendor) :] sum_data(combined_item, item, wrap_none=wrap_none) combined_item["tier_level"][sub_dsp] = sub_dsp_streams return [ date_item for isrc_item in results.values() for market_item in isrc_item.values() for date_item in market_item.values() ] async def get_videos(*args, **kwargs) -> list: """Emulate get videos Delphi endpoint.""" kwargs = kwargs or {} kwargs.update(zip(streams_args_spec, args)) if "isrc_list" not in kwargs: return [] views_count = kwargs.get("views", 0) data = read_json_from_file("videos_by_isrc") results = [] for item in data: if item["views"] < views_count: continue results.append(item) return results def filter_tiktok_result(data: dict, fields: Iterable[str], level: int): """Filter tiktok track analytics result. Args: data: Data to filter in. fields: Fields to have in result. level: How deep inside data we need to dig into to find items to filter in. """ if level > 0: for value in data.values(): filter_tiktok_result(value, fields, level - 1) else: for key in list(data.keys()): if key not in fields: del data[key] async def get_tiktok_tracks_analytics(*args, **kwargs) -> dict: """Emulate get /tiktok/tracks/analytics Delphi endpoint.""" min_date = date(2020, 9, 2) max_date = date(2020, 10, 17) kwargs = kwargs or {} kwargs.update(zip(tiktok_tracks_analytics_args_spec, args)) start_date = _get_date_from_kwargs("start_date", kwargs) if start_date < min_date: start_date = min_date end_date = _get_date_from_kwargs("end_date", kwargs) if end_date > max_date: end_date = max_date result = read_json_from_file("tiktok_tracks_analytics") result["min_date"] = start_date.isoformat() result["max_date"] = end_date.isoformat() days_count = (end_date - start_date).days + 1 result["days_count"] = days_count range_persent = days_count / ((max_date - min_date).days + 1) start = (start_date - min_date).days end = (end_date - min_date).days + 1 content_type_list = kwargs.get("content_type_list") metrics_list = kwargs.get("metrics_list") country_code_list = kwargs.get("country_code_list") breakdowns_list = kwargs.get("breakdowns_list") or [TiktokTrackAnalyticsBreakdown.DAILY] for breakdown_name, breakdown_values in result["breakdowns"].items(): total = None if metrics_list: filter_tiktok_result( breakdown_values, metrics_list, TiktokTrackAnalyticsBreakdown.METRIC_DEPTH[breakdown_name] ) if breakdown_name == TiktokTrackAnalyticsBreakdown.DAILY: if TiktokTrackAnalyticsBreakdown.TOTALS in breakdowns_list and not country_code_list: total = {} for field, values in list(breakdown_values.items()): if metrics_list and field not in metrics_list: del breakdown_values[field] else: breakdown_values[field] = values[start:end] if total is not None: total[field] = sum(i or 0 for i in values) elif breakdown_name == TiktokTrackAnalyticsBreakdown.CONTENT_TYPE_COUNTRY: if TiktokTrackAnalyticsBreakdown.TOTALS in breakdowns_list and country_code_list: total = defaultdict(lambda: defaultdict(int)) for content_type_name, content_type_values in list(breakdown_values.items()): if content_type_list and content_type_name not in content_type_list: del breakdown_values[content_type_name] continue for country_code_name, country_code_values in list(content_type_values.items()): if (country_code_list and country_code_name not in country_code_list) or ( not country_code_list and country_code_name != DELPHI_GLOBAL_MARKET ): del content_type_values[country_code_name] continue for field, values in list(country_code_values.items()): if metrics_list and field not in metrics_list: del country_code_values[field] continue country_code_values[field] = values[start:end] if total is not None: total[country_code_name][field] = sum(i or 0 for i in values) elif breakdown_name == TiktokTrackAnalyticsBreakdown.COUNTRY_TOTALS: for country_code_name, country_code_values in list(breakdown_values.items()): if (country_code_list and country_code_name not in country_code_list) or ( not country_code_list and country_code_name != DELPHI_GLOBAL_MARKET ): del breakdown_values[country_code_name] continue for field in country_code_values.keys(): country_code_values[field] = math.ceil(country_code_values[field] * range_persent) if total is not None: result["breakdowns"][TiktokTrackAnalyticsBreakdown.TOTALS] = total result["breakdowns"] = {key: value for key, value in result["breakdowns"].items() if key in breakdowns_list} return result def get_regions(*args, **kwargs) -> dict: """Emulate get regions Delphi endpoint.""" kwargs = kwargs or {} kwargs.update(zip(regions_args_spec, args)) sort_by, sort_order, limit, offset = ( kwargs.get("sort_by", "country_code"), kwargs.get("sort_order", "asc"), kwargs.get("limit", 1000), kwargs.get("offset", 0), ) data = read_json_from_file("regions") data["items"] = list(sorted(data["items"], key=lambda x: x[sort_by], reverse=(sort_order == "desc")))[ offset : offset + limit ] return data def get_tracks(*args, **kwargs) -> list: """Emulate get tracks Delphi endpoint.""" mapping = { "USSM19902989": ["USSM19902989", "USSM19902990", "USSM19902991"], "USSM19902990": ["USSM19902990", "USSM19902991"], "USSM19902991": ["USSM19902990", "USSM19902991"], "USSM11122201": ["USSM11122201", "USSM11122202"], "USSM11122202": ["USSM11122201", "USSM11122202"], } isrc_list = kwargs["isrc"] return [{"isrc": isrc, "related_isrcs": list(mapping.get(isrc, [isrc]))} for isrc in isrc_list] @pytest.fixture def mocked_streams(): with patch.object(delphi_client.DelphiClient, "get_streams") as mocked_get_streams: mocked_get_streams.side_effect = get_streams yield mocked_get_streams @pytest.fixture def mocked_first_dates(): with patch.object(delphi_client.DelphiClient, "get_first_stream_dates") as mocked_get_first_stream_dates: mocked_get_first_stream_dates.side_effect = get_first_streams_dates yield mocked_get_first_stream_dates @pytest.fixture def mocked_latest_dates(): with patch.object(delphi_client.DelphiClient, "get_latest_stream_dates") as mocked_get_latest_stream_dates: mocked_get_latest_stream_dates.return_value = latest_dates_data yield mocked_get_latest_stream_dates @pytest.fixture def mocked_videos(): with patch.object(delphi_client.DelphiClient, "get_videos") as mocked_get_videos: mocked_get_videos.side_effect = get_videos yield mocked_get_videos @pytest.fixture def mocked_video_analytics(): with patch.object(delphi_client.DelphiClient, "get_video_analytics") as mocked_get_video_analytics: mocked_get_video_analytics.side_effect = get_video_analytics yield mocked_get_video_analytics @pytest.fixture def mocked_tiktok_tracks_analytics(): with patch.object(delphi_client.DelphiClient, "_get_tiktok_tracks_analytics") as mocked_get_tiktok_analytics: mocked_get_tiktok_analytics.side_effect = get_tiktok_tracks_analytics yield mocked_get_tiktok_analytics @pytest.fixture def mocked_regions(): with patch.object(delphi_client.DelphiClient, "get_regions") as mocked_get_regions: mocked_get_regions.side_effect = get_regions yield mocked_get_regions @pytest.fixture def mocked_tracks(): with patch.object(delphi_client.DelphiClient, "get_tracks") as mocked_get_tracks: mocked_get_tracks.side_effect = get_tracks yield mocked_get_tracks async def get_charts(*args, **kwargs) -> dict: """Emulate get charts Delphi endpoint.""" kwargs = kwargs or {} kwargs.update(zip(charts_args_spec, args)) dsp = kwargs.get("dsp") sort_by = kwargs.get("sort_by") sort_order = kwargs.get("sort_order") filename = f"{'charts_summary'}" data = read_json_from_file(filename) if sort_by: if sort_order == "asc": data = sorted(data, key=lambda k: k["rank"]) elif sort_order == "desc": data = sorted(data, key=lambda k: k["rank"], reverse=True) if dsp: dsp = dsp.lower().split(",") result = [] for i in data: if i["dsp"]["name"].lower() not in dsp: continue result.append(i) return {"count": len(result), "items": result} return {"count": len(data), "items": data} @pytest.fixture def mocked_charts(): with patch.object(delphi_client.DelphiClient, "get_charts") as mocked_get_charts: mocked_get_charts.side_effect = get_charts yield mocked_get_charts async def get_charts_analytics(kwargs) -> dict: """Emulate get charts analytics Delphi endpoint.""" isrc_list, start_date, end_date = ( kwargs.get("isrc"), str_to_date(kwargs["start_date"]), str_to_date(kwargs["end_date"]), ) if isinstance(isrc_list, str): isrc_list = [isrc_list] data = read_json_from_file("spotify_charts_analytics") for item in list(data): if isrc_list and item["isrc"] not in isrc_list: data.remove(item) continue del item["isrc"] min_date, max_date, positions = ( str_to_date(item["min_date"]), str_to_date(item["max_date"]), item["metrics"]["positions"], ) if min_date < start_date: positions = positions[(start_date - min_date).days :] item["min_date"] = None if start_date > max_date else start_date.isoformat() if max_date > end_date: positions = positions[: -(max_date - end_date).days] item["max_date"] = None if end_date < min_date else end_date.isoformat() if not positions: data.remove(item) continue item["metrics"]["positions"] = positions item["min_position"] = nullable_min(*positions) item["max_position"] = nullable_max(*positions) return {"count": len(data), "items": data} async def get_tracks_charts(kwargs) -> dict: """Emulate get tracks charts Delphi endpoint.""" isrc_list, chart_type, chart_breakdown = kwargs.get("isrc"), kwargs.get("chart_type"), kwargs.get("chart_breakdown") if isrc_list and isinstance(isrc_list, str): isrc_list = [isrc_list] data = read_json_from_file("tracks_charts") for item in list(data): if ( (isrc_list and item["public_meta"]["isrc"] not in isrc_list) or (chart_type and item["chart_meta"]["type"] != chart_type) or (chart_breakdown and item["chart_meta"]["breakdown"] != chart_breakdown) ): data.remove(item) continue return {"count": len(data), "items": data} async def get(*args, **kwargs): """Mock Delphi client get method.""" kwargs = kwargs or {} kwargs.update(zip(get_args_spec, args)) url, params = kwargs["url"], kwargs["params"] parsed_params = {} for key, value in params: if key in parsed_params: parsed_value = parsed_params[key] if isinstance(parsed_value, list): parsed_value.append(value) else: parsed_params[key] = [parsed_value, value] else: parsed_params[key] = value if url in ("spotify/charts/analytics", "apple-music/charts/analytics"): return await get_charts_analytics(parsed_params) elif url in ("spotify/tracks/charts", "apple-music/tracks/charts"): return await get_tracks_charts(parsed_params) raise NotImplementedError() @pytest.fixture def mocked_get(): with patch.object(delphi_client.DelphiClient, "_get") as mocked_get: mocked_get.side_effect = get yield mocked_get