import asyncio import logging from apollo_utils.core.utils.dispatchers.dsp_dispatcher import dsp_dispatch from collections import defaultdict from datetime import date, timedelta from dateutil import relativedelta from typing import Any, Dict, List, Optional from server import config from server.cache.utils import cached from server.client import services from server.constants import DEFAULT_MAX_PAGINATION_LIMIT, DSP from server.constants.playlists.v0.misc import StreamsType, StreamsTypeRequest from server.utils.playlists.v0.apple.playlist_images import get_apple_playlist_images from server.utils.playlists.v0.spotify.playlist_images import get_spotify_playlists_images logger = logging.getLogger("app") @cached(ttl=config.GET_PLAYLISTS_BY_TRACK_DATA_TTL) async def get_playlists_data( dsp: DSP, isrc: str, streams_type: StreamsTypeRequest, markets_list: Optional[List[str]] = None, start_date: Optional[date] = None, end_date: Optional[date] = None, search: Optional[str] = None, category_id: Optional[int] = None, recent_adds_only: bool = False, keep_zero_streams: bool = True, ): """Get playlists data from Apollo API and DSP API. Args: dsp: Apple. isrc: track isrc to playlists by. streams_type: type of requested streams. markets_list: markets to filter playlists by. start_date: min date of track in playlists presence to filter by. end_date: max date of track in playlists presence to filter by. search: string to filter playlists name by. category_id: playlist category_id to filter_by. recent_adds_only: flag to return only playlists that track was recent added to. keep_zero_streams: flag to keep playlists with 0 streams in a result or not. Returns: Playlists, Streams type """ if not start_date: start_date = ( end_date - relativedelta.relativedelta(days=7) if end_date else date.today() - relativedelta.relativedelta(days=7) ) if not end_date: end_date = start_date + relativedelta.relativedelta(days=7) if start_date else date.today() start_date, end_date = start_date.isoformat(), end_date.isoformat() results = [] if dsp in (DSP.SPOTIFY, DSP.APPLE): results = await services.apollo.get_track_playlists( dsp=dsp, isrc=isrc, market=markets_list, recent_adds_only=recent_adds_only, include=["followers", "owner"], limit=DEFAULT_MAX_PAGINATION_LIMIT, search=search, category_id=category_id, ) results = results["items"] elif dsp == DSP.AMAZON: dirty_results = await services.dsp.get_track_positions_playlists( isrc=isrc, streams_country_code=markets_list, country_code=markets_list, dsp=dsp.value, include="playlists", ) for row in dirty_results: if not row.get("playlist_id"): logger.warning(f"Got not valid row for isrc {isrc}.\n") continue row["id"] = row["playlist"]["amazon_playlist_id"] row["country_code"] = row["playlist"]["country_code"] results.append(row) items = defaultdict(set) for i in results: items[i["id"]].add(i["country_code"] or "") items = {k: list(v) for k, v in items.items()} streams_request_type_to_getter_and_result_type = { StreamsTypeRequest.TRACK: ( [ services.dsp.get_track_or_playlists_streams( vendor=dsp.value, isrc=isrc, items=items, start_date=start_date, end_date=end_date ), ], StreamsType.TRACK, ), StreamsTypeRequest.PLAYLIST: ( [ services.dsp.get_track_or_playlists_streams( vendor=dsp.value, items=items, start_date=start_date, end_date=end_date ), ], StreamsType.PLAYLIST, ), } streams_request_type_to_getter_and_result_type[StreamsTypeRequest.PRIORITY] = ( streams_request_type_to_getter_and_result_type[StreamsTypeRequest.TRACK][0] + streams_request_type_to_getter_and_result_type[StreamsTypeRequest.PLAYLIST][0], None, ) tasks, result_streams_type = streams_request_type_to_getter_and_result_type[streams_type] streams_result = await asyncio.gather(*tasks) if not result_streams_type: track_in_playlist_streams_result, playlist_streams_result = streams_result if any([i["current_streams"] for i in track_in_playlist_streams_result]): result_streams_type, streams_result = StreamsType.TRACK, [track_in_playlist_streams_result] else: result_streams_type, streams_result = StreamsType.PLAYLIST, [playlist_streams_result] items = add_streams(results, streams_result[0], keep_zero_streams) return items, result_streams_type.value def add_streams(results: List[Dict[str, Any]], dsp_results: List[Dict[str, Any]], keep_zero_streams=True): """Add dps streams numbers data.""" key = "{id}_{country_code}" dsp_dict = {key.format(**item): item for item in dsp_results} for item in results: item["prev_streams"] = dsp_dict.get(key.format(**item), {}).get("prev_streams") item["current_streams"] = dsp_dict.get(key.format(**item), {}).get("current_streams") if not keep_zero_streams: results = [i for i in results if i["current_streams"]] return results def sum_periods(dates: List[List[date]]) -> List[List[date]]: """Concat overlapping and consecutive date intervals. Args: dates (List[List[date]]): List of date pairs what are interval starts and ends. Returns: List[List[date]]: Aggregated list of date pairs. """ if not dates: return dates dates = sorted(dates, key=lambda x: (x[0], x[1])) prev_date = dates[0] results = [prev_date] for i in range(1, len(dates)): current_date = dates[i] if date.fromisoformat(prev_date[1]) >= (date.fromisoformat(current_date[0]) - timedelta(days=2)): if current_date[1] > prev_date[1]: prev_date[1] = current_date[1] else: prev_date = current_date results.append(prev_date) return results def get_playlists_previous_top_response( dsp: DSP, history_mapping: dict[str, dict[str, list[list[date]]]], id_to_name_map: dict[str, str], streams_map: dict[str, int], tracks_mapping: dict[str, dict], playlist_id_to_image_url_map: dict[str, str], track_first_dates: dict[str, str], limit: int, ): result = [] for playlist_id, playlist_tracks in history_mapping.items(): playlist = { "id": playlist_id, "name": id_to_name_map.get(playlist_id), "streams": int(streams_map.get(playlist_id, 0)), "type": dsp.value, "image_url": playlist_id_to_image_url_map.get(playlist_id, ""), "tracks": [], } result.append(playlist) for track_id, history_dates in playlist_tracks.items(): track_data = tracks_mapping[str(track_id)] track = { "id": track_id, "name": track_data["name"], "isrc": track_data["isrc"], "first_date": track_first_dates[track_data["isrc"]] if track_data["isrc"] in track_first_dates else False, "periods": sum_periods(history_dates), } playlist["tracks"].append(track) result = sorted(result, key=lambda x: (-x["streams"])) if limit: result = result[:limit] return {"playlists": result} @dsp_dispatch((DSP.SPOTIFY, get_spotify_playlists_images), (DSP.APPLE, get_apple_playlist_images)) async def get_dsp_playlists_images( playlist_ids: List[int], market: str = None, image_size: int = None, dsp: DSP = None ): raise NotImplemented("Unknown dsp")