import asyncio from aiohttp.web import Request, Response, json_response from datetime import date, timedelta from typing import List, Optional from server.client.clients.delphi_client import DelphiClient from server.constants.delphi.streams.group_by import StreamsGroupBy from server.constants.delphi.streams.include import StreamsInclude from server.legacy.analytics import constants as analytics_constants from server.legacy.analytics import utils as analytics_utils from server.legacy.core import constants, utils from server.legacy.delphi import constants as delphi_constants from server.utils.delphi.converters import str_to_date async def get_streams_monitoring_data(request: Request, dsp_list: List[str]) -> Response: api = request.app["delphi_java_api"] data = request["querystring"] tasks = [ api.get_streams( start_date=data["start_date"], end_date=data["end_date"], isrc_list=[data["isrc"]], dsp_list=[vendor], country_code_list=[constants.DELPHI_GLOBAL_MARKET], include=[StreamsInclude.SOURCES] if vendor == constants.SPOTIFY else None, ) for vendor in dsp_list ] result = {} responses = await asyncio.gather(*tasks) for vendor, response in zip(dsp_list, responses): result.update({vendor: analytics_utils.process_vendor_streaming_data(vendor, response)}) return json_response(result) async def get_date_range_streams(request: Request, dsp_list: List[str], streams_field: str = "results") -> Response: api = request.app["delphi_java_api"] data = request["querystring"] response = await api.get_streams( start_date=data["start_date"], end_date=data["end_date"], isrc_list=[data["isrc"]], dsp_list=dsp_list, country_code_list=[data["market"]], ) streams_mapping = {item["dsp"]: item.get("streams", 0) for item in response} results = {streams_field: sum(streams_mapping.values())} if data.get("per_vendors"): results["vendors"] = streams_mapping return json_response(results) async def get_stream_graph( request: Request, dsp_list: List[str], with_sources: bool = False, show_total: bool = False, wrap_single_vendor: bool = False, add_market: bool = False, default_none: bool = False, base_mapping: dict = analytics_constants.DETAILED_STREAMS_MONITORING_BASE, source_mapping: dict = analytics_constants.DETAILED_STREAMS_MONITORING_SOURCE_V1, ) -> Response: """Get streams per day as coordinated. Args: request: Request object. dsp_list: Vendor list. with_sources: Only streams or all sources. show_total: Sum all vendors streams. wrap_single_vendor: When vendor param is set wrap in coordinates dict. add_market: Add market to response. default_none: Display empty Y value as None - if True, or as zero - if False. base_mapping: Base mapping object for total streams data. source_mapping: Base mapping object for vendors streams sources data. Returns: Streams per day as response object. """ api = request.app["delphi_java_api"] data = request["querystring"] start_date = data["start_date"] end_date = data["end_date"] market = data["market"] vendor = data.get("vendor") dsp_list = [vendor] if vendor else dsp_list response = await api.get_streams( start_date=start_date, end_date=end_date, isrc_list=[data["isrc"]], dsp_list=dsp_list, country_code_list=[market], group_by=[StreamsGroupBy.DATE], include=[StreamsInclude.SOURCES] if with_sources else None if with_sources else None, ) results = analytics_utils.process_stream_graph_data( response, start_date, end_date, dsp_list, show_total=show_total, per_vendors=data.get("per_vendors"), base_mapping=base_mapping if with_sources else None, source_mapping=source_mapping if with_sources else None, default_none=default_none or with_sources, ) if vendor and not show_total: results = results[vendor] if wrap_single_vendor: results = {"coordinates": results} if add_market and (not vendor or wrap_single_vendor or show_total): results["market"] = utils.convert_market(market, constants.GLOBAL_MARKET) return json_response(results) async def get_total_streams_count( request: Request, dsp_list: List[str], first_stream_dates: bool, global_field: str = "global", market_field: str = "market", ) -> Response: api = request.app["delphi_java_api"] data = request["querystring"] isrc_list = data.get("isrc_list") or [data["isrc"]] market = data.get("market") end_date = await analytics_utils.get_last_date_value(api=api, data=data, key="end_date") tasks = [ api.get_streams( start_date=delphi_constants.START_DATE_MIN, end_date=end_date, isrc_list=isrc_list, dsp_list=dsp_list, country_code_list=[constants.DELPHI_GLOBAL_MARKET] + ([market] if market else []), ) ] if first_stream_dates: tasks.append(api.get_first_stream_dates(isrc_list=isrc_list, dsp_list=dsp_list)) responses = await asyncio.gather(*tasks) streams_data = responses[0] total_streams, market_streams = analytics_utils.process_total_streams_count(streams_data, True) result = {global_field: sum(total_streams.values())} if market: result.update({market_field: sum(market_streams.values())}) if data.get("per_vendors"): result_per_vendor = {} for vendor in dsp_list: result_per_vendor[vendor] = {global_field: total_streams[vendor]} if market: result_per_vendor[vendor][market_field] = market_streams[vendor] result["vendors"] = result_per_vendor if first_stream_dates: result["first_stream_date"] = analytics_utils.process_first_stream_dates(responses[1], True) return json_response(result) async def get_daily_and_weekly_based_streams_data( request: Request, dsp_list: List[str], with_dates: bool, days_and_chart_weeks: bool = True ) -> Response: """Calculate chart week and daily based data fo track by it ISRC in global and specific market.""" api = request.app["delphi_java_api"] data = request["querystring"] market, isrc = data["market"], data["isrc"] last_date = await analytics_utils.get_last_date_value(api=api, data=data, key="last_date") prev_date = last_date - timedelta(weeks=1) current_week, prev_week = ( analytics_utils.get_chart_weeks_ranges(last_date) if days_and_chart_weeks else analytics_utils.get_simple_weeks_ranges(last_date) ) rolling_week = analytics_utils.get_rolling_week() if days_and_chart_weeks else None response = await api.get_streams( start_date=prev_week.start, end_date=rolling_week.end if days_and_chart_weeks else current_week.end, isrc_list=[isrc], dsp_list=dsp_list, country_code_list=[constants.DELPHI_GLOBAL_MARKET] + ([market] if market else []), group_by=[StreamsGroupBy.DATE], ) result = analytics_utils.get_daily_and_weekly_streams( response, current_week, prev_week, rolling_week, last_date, prev_date, dsp_list, market, data.get("per_vendors"), days_and_chart_weeks, ) if with_dates: result_dates = { "current_week": current_week.to_dict(isoformat=True), "prev_week": prev_week.to_dict(isoformat=True), } if days_and_chart_weeks: result_dates.update( { "rolling_week": rolling_week.to_dict(isoformat=True), "current_day": last_date.isoformat(), "prev_day": prev_date.isoformat(), } ) result["dates"] = result_dates return json_response(result) async def get_tracks_streams(request: Request) -> Response: api = request.app["delphi_java_api"] data = request["querystring"] vendor = data.get("vendor") dsp_list = data.get("vendors") or ([vendor] if vendor else constants.VENDORS) response = await api.get_streams( start_date=data["start_date"], end_date=data["end_date"], isrc_list=data["isrc_list"], dsp_list=dsp_list, country_code_list=[data["market"]], ) streams_sum, streams_per_vendor = analytics_utils.process_streams_count(response) results = {"streams": streams_sum} if data.get("per_vendors"): results["vendors"] = streams_per_vendor return json_response(results) async def get_shazam_charts_positions_and_summary( api: DelphiClient, chart_id: str, chart_date: date, is_sony: Optional[bool] ) -> List[dict]: """Get shazam charts positions and summary data together. Args: api: Delphi client. chart_id: Chart ID like "shazam_us" or "shazam_us_7060" chart_date: Chart date. is_sony: Is Sony filter. Returns: Positions and summary shazam charts data. """ result = await api.get_shazam_charts_track_positions( start_date=chart_date.isoformat(), end_date=chart_date.isoformat(), chart_id=chart_id, include="track", is_sony=is_sony, ) isrc_list = [i["isrc"] for i in result if i.get("isrc")] summary_data = await api.get_shazam_charts_track_positions_summary(chart_id=chart_id, isrc=isrc_list) summary_mapping = {i["isrc"]: i["summary"] for i in summary_data} for item in result: isrc = item.get("isrc") if isrc and isrc in summary_mapping: item["summary"] = summary_mapping[isrc] return result async def get_shazam_charts_position(request: Request) -> dict: """Get shazam charts positions data. Args: request: Web request. Returns: Shazam charts positions data. """ api, data = request.app["delphi_java_api"], request["querystring"] start_date, end_date, isrc, market, city = ( data["start_date"], data["end_date"], data["isrc"], data["market"], data["city"], ) chart_id = analytics_utils.get_shazam_chart_id(market, city) tasks = [ api.get_shazam_charts_track_positions(start_date=start_date, end_date=end_date, isrc=isrc, chart_id=chart_id) ] if city: tasks.append(api.get_shazam_cities(city_id=city)) responses = await asyncio.gather(*tasks) positions = responses[0] city_data = responses[-1][0] if city and responses[-1] else None min_position, max_position, min_date, max_date = 1000, 0, end_date, start_date for item in positions: position, chart_date = item["position"]["current"], str_to_date(item["position"]["date"]) if position < min_position: min_position = position if position > max_position: max_position = position if chart_date < min_date: min_date = chart_date if chart_date > max_date: max_date = chart_date return { "isrc": isrc, "start_date": min_date, "end_date": max_date, "days": (max_date - min_date).days + 1, "city": city_data, "market": market, "position": {"peak": min_position if positions else None, "min": max_position if positions else None}, "values": positions, }