import asyncio import logging from apollo_utils.core.utils.ext_enum import ALL from apollo_utils.service.exceptions import APIInvalidResponse from datetime import date from typing import List, Optional, Tuple from server.client.base.client import ApiKeyClient from server.client.utils import prepare_list_arg, str_to_date from server.constants import DSP, DspApiPrefix from server.constants.charts import ChartStatus, SourceType, SourceTypeDataHealthChartUpdateField, TrackStateIncludes from server.domains.charts.misc import get_chart_id from server.utils.parallel import DictResult, ListResult, request_in_chunks logger = logging.getLogger(__name__) def merge_tracks_brands_tagging_chunk_result(result: dict, chunk_result: dict) -> dict: keys = ("items", "unknown_upcs", "unknown_track_ids") for key in keys: result[key] = (result.get(key) or []) + ((chunk_result or {}).get(key) or []) return result class DpsApiDelphiClient(ApiKeyClient): async def get_artists(self, **params): result = await self.send_request("api/delphi/artists", params=params) return result["items"] @staticmethod def _prepare_tracks_charts(dsp: DSP, params: dict) -> dict: if "metrics_dimension" not in params: params["metrics_dimension"] = "isrc" for name in ("chart_breakdown", "chart_type"): if params.get(name) and (dsp.value == DSP.APPLE.value or len(params[name]) > 1): del params[name] return params async def get_chart_positions_dates(self, **params): return await self.send_request("api/delphi/data-health/chart-positions", params=params) async def get_charts(self, **params): return await self.send_request("api/delphi/charts", params=params) async def get_charts_additions(self, dsp: DSP, **params): if "metrics_dimension" not in params: params["metrics_dimension"] = "isrc" result = await self.send_request( f"api/delphi/{DspApiPrefix[dsp.name].value}/charts/tracks/additions", params=params ) return result["items"] async def get_charts_analytics(self, **params): if "metrics" not in params: params["metrics"] = "positions" return await self.send_request("api/delphi/charts/analytics", params=params) async def get_charts_data_health_status(self, dsp: DSP, **params): return await self.send_request( f"api/delphi/{DspApiPrefix[dsp.name].value}/charts/data-health/status", params=params ) async def get_charts_latest_date( self, dsp: DSP, chart_breakdown: str, chart_type: str, country_code: str, source_type: Optional[SourceType] = None, ) -> date: data_health_chart_updated_key = ( SourceTypeDataHealthChartUpdateField[dsp.value] if source_type is not None else "max_date" ) result = await self.get_charts_latest_date_bulk( dsp=dsp, chart_breakdown=[chart_breakdown], chart_type=[chart_type], country_code=[country_code], data_health_chart_updated_key=data_health_chart_updated_key, ) return list(result.values())[0] if result else None async def get_charts_latest_date_bulk( self, dsp: DSP, chart_breakdown: List[str], chart_type: List[str], country_code: Optional[List[str]] = None, data_health_chart_updated_key: str = "max_date", ): if chart_breakdown: chart_breakdown = None if len(chart_breakdown) > 1 else chart_breakdown[0] if chart_type: chart_type = None if len(chart_type) > 1 else chart_type[0] params = {} if dsp.value == DSP.SPOTIFY.value: params = {"chart_breakdown": chart_breakdown, "chart_type": chart_type} result_field = "breakdown_chart_type_country_code" else: result_field = "dsp_chart_type_country_code" result = await self.get_charts_data_health_status( dsp=dsp, status=[ChartStatus.COMPLETE.value, ChartStatus.COMPLETE_VOLATILE.value], country_code=country_code, **params, ) return { get_chart_id(country_code_name, type_name, breakdown_name, dsp): str_to_date( country_dates[data_health_chart_updated_key] ) for breakdown_name, type_items in result[result_field].items() for type_name, country_items in type_items.items() for country_code_name, country_dates in country_items.items() } async def get_charts_removals(self, dsp: DSP, **params): if "metrics_dimension" not in params: params["metrics_dimension"] = "isrc" result = await self.send_request( f"api/delphi/{DspApiPrefix[dsp.name].value}/charts/tracks/removals", params=params ) return result["items"] async def get_charts_tracks(self, dsp: DSP, **params): if "metrics_dimension" not in params: params["metrics_dimension"] = "isrc" result = await self.send_request(f"api/delphi/{DspApiPrefix[dsp.name].value}/charts/tracks", params=params) return result["items"] async def get_dsp_charts(self, dsp: DSP, **params): if dsp.value == DSP.APPLE.value: for name in ("breakdown", "type"): if name in params: del params[name] result = await self.send_request(f"api/delphi/{DspApiPrefix[dsp.name].value}/charts", params=params) return result["items"] async def get_juno_spotify_streams_latest(self, **params): result = await self.send_request("api/delphi/juno/spotify/streams/latest", params=params) return result["items"] async def get_latest_streams_date(self, dsp: Optional[DSP] = None, use_min: bool = False) -> Optional[date]: result = await self.send_request("api/delphi/data-health/completeness") if dsp is None: if use_min: min_date = None for item in result: dsp_date = str_to_date(item["updated_date"]) if min_date is None or dsp_date < min_date: min_date = dsp_date return min_date return result for item in result: if item["dsp"] == dsp.value: return str_to_date(item["updated_date"]) async def get_public_current_playlists_info(self, playlist_id, allowed_status_codes: Tuple = None, **params): try: return await self.send_request(f"api/delphi/public/playlists/current/{playlist_id}", params=params) except APIInvalidResponse as vendor_error: if allowed_status_codes and vendor_error.original_status_code in allowed_status_codes: return vendor_error.original_response else: raise vendor_error @request_in_chunks(chunk_size=150, items_key="playlist_id", result_type=ListResult) async def get_public_playlists(self, **params): result = await self.send_request("api/delphi/public/playlists", params=params) return result.get("items", []) async def get_public_playlists_info(self, playlist_id, allowed_status_codes: Tuple = None, **params): try: return await self.send_request(f"api/delphi/public/playlists/{playlist_id}", params=params) except APIInvalidResponse as vendor_error: if allowed_status_codes and vendor_error.original_status_code in allowed_status_codes: return vendor_error.original_response else: raise vendor_error async def get_public_track_positions_playlists(self, **params): return await self.send_request("api/delphi/public/track-positions/playlists", params=params) async def get_public_track_positions_playlists_current_tracklist(self, **params): return await self.send_request("api/delphi/public/track-positions/playlists/current-tracklist", params=params) async def get_public_track_positions_playlists_dates(self, **params): return await self.send_request("api/delphi/public/track-positions/playlists/dates", params=params) async def get_public_track_positions_previous_playlists(self, **params): return await self.send_request("api/delphi/public/track-positions/previous/playlists", params=params) async def get_regions(self, **params): result = await self.send_request("api/delphi/regions", params=params) return result["items"] async def get_search(self, **params): result = await self.send_request("api/delphi/search", params=params) return result["items"] @prepare_list_arg("isrc") @request_in_chunks(chunk_size=150, items_key="playlist_id", result_type=ListResult) async def get_streams(self, **params): return await self.send_request("api/delphi/streams", params=params) async def get_streams_country_code_rank(self, **params): return await self.send_request("api/delphi/streams/country-code/rank", params=params) async def get_streams_insights(self, **params): return await self.send_request("api/delphi/streams/insights", params=params) @prepare_list_arg("isrc") @request_in_chunks(chunk_size=100, items_key="isrc", result_type=ListResult) async def get_streams_isrc_chunks(self, **params): return await self.send_request("api/delphi/streams", params=params) @prepare_list_arg("isrc") async def get_tiktok_top_tracks(self, **params) -> dict: """Get tiktok top tracks (charts) from Delphi API. Returns: Tiktok top tracks charts. """ result = await self.send_request("api/delphi/tiktok/top/tracks", params=params) return result["items"] @prepare_list_arg("isrc") async def get_tiktok_top_tracks_analytics(self, **params) -> dict: """Get tiktok top tracks analytics (day by day positions) from Delphi API. Returns: Tiktok top tracks positions for dates range. """ return await self.send_request("api/delphi/tiktok/top/tracks/analytics", params=params) async def get_tiktok_tracks_analytics(self, **params): return await self.send_request("api/delphi/tiktok/tracks/analytics", params=params) async def get_track_playlists(self, **params): return await self.send_request("api/delphi/track-positions/playlists", params=params) async def get_track_positions_playlists(self, **params): result = await self.send_request("api/delphi/track-positions/playlists", params=params) return result.get("items", []) async def get_tracks(self, **params): result = await self.send_request("api/delphi/tracks", params=params) return result["items"] @request_in_chunks( chunk_size=100, result_type=DictResult, items_key=("track_id", "upc"), update_func=merge_tracks_brands_tagging_chunk_result, ) async def get_tracks_brands_tagging(self, **params) -> dict: # keep in mind that 'worldwide' works as 'any' market, in fact there is no 'worldwide' market for upc, # so delphi returns a relation for track_id/upc and a distributor tuple in 'worldwide' if there is # a relation for this tuple in any market return await self.send_request("api/delphi/tracks/brands-tagging", params=params) @prepare_list_arg("isrc") async def get_tracks_charts(self, dsp: DSP, **params): params = self._prepare_tracks_charts(dsp, params) result = await self.send_request(f"api/delphi/{DspApiPrefix[dsp.name].value}/tracks/charts", params=params) return result["items"] @prepare_list_arg("isrc") async def get_tracks_charts_lifetime(self, dsp: DSP, **params): params = self._prepare_tracks_charts(dsp, params) result = await self.send_request( f"api/delphi/{DspApiPrefix[dsp.name].value}/tracks/charts/lifetime", params=params ) return result["items"] async def get_tracks_charts_lifetime_with_metrics(self, dsp: DSP, **params): if ( "track_state_includes" in params and TrackStateIncludes.CHART_META.value not in params["track_state_includes"] ): params["track_state_includes"].append(TrackStateIncludes.CHART_META.value) tasks = [ self.get_tracks_charts( dsp=dsp, track_state_includes=[TrackStateIncludes.CHART_META.value, TrackStateIncludes.METRICS.value], **params, ), self.get_tracks_charts_lifetime(dsp=dsp, **params), ] metrics, result = await asyncio.gather(*tasks) metrics = {i["chart_meta"]["chart_id"]: i["metrics"] for i in metrics} for item in result: chart_id = item["chart_meta"]["chart_id"] if chart_id in metrics: item["metrics"] = metrics[chart_id] return result @prepare_list_arg("isrc") async def get_tracks_in_chart(self, **params): response = await self.get_tracks_charts_lifetime( dsp=DSP.SPOTIFY, track_state_includes=[TrackStateIncludes.PUBLIC_META.value], **params ) return [i["public_meta"]["isrc"] for i in response] @prepare_list_arg("video_id") async def get_video_positions_charts(self, **params): return await self.send_request("api/delphi/video-positions/charts", params=params) @prepare_list_arg("video_id") async def get_video_positions_charts_summary(self, **params): return await self.send_request("api/delphi/video-positions/charts/summary", params=params)