"""Logic for retrieving streams timeseries by store.""" import datetime from typing import Any, Mapping import oto.response from ddtrace import tracer from analytics.constants import cache from analytics.constants.parameters import ALL_TIME from analytics.logic import data_availability from analytics.queries.sound_recording_streams import StreamsByStore from analytics.schemas.streams_by_store import StreamsByStoreSchema from analytics.utils import store_availability from analytics.utils import streams as streams_utils from analytics.utils.cache import cache_in_redis from analytics.validation.schema import schema_dump @cache_in_redis(ttl=cache.ONE_DAY) @tracer.wrap(name="get_streams_by_store") def get_streams_by_store( query_params: Mapping[str, Any], permissions: Mapping[str, Any], ): """Return streams by store timeseries for ISRC. Args: query_params: Dict with isrc, distributors, country_ids, store_ids, start_date, end_date. permissions: Dict with permission_* keys. Returns: oto.response.Response with streams by store payload. """ isrc = query_params["isrc"] response_body = {"isrc": isrc, "stores": []} start_date = query_params.get("start_date") end_date = query_params.get("end_date") if not (start_date and end_date): start_date, end_date = data_availability.get_date_range( data_availability.HIGHWATERMARK_DATE, days=7 ) all_time = start_date == ALL_TIME store_ids = query_params.get("store_ids", []) if not store_ids: store_ids = store_availability.get_store_ids() else: store_ids = sorted( list(set(store_ids).intersection(store_availability.get_store_ids())) ) if not store_ids: schema = StreamsByStoreSchema() return oto.response.Response(schema_dump(schema, response_body)) end_date_str = ( end_date.strftime("%Y-%m-%d") if isinstance(end_date, datetime.date) else end_date ) query_input = { **permissions, "isrc": isrc, "store_ids": store_ids, "distributors": query_params["distributors"], "end_date": end_date_str, "country_ids": query_params.get("country_ids", []), "all_time": all_time, "transfer_product_ownership_enabled": query_params.get( "transfer_product_ownership_enabled", False ), } if not all_time: start_date_str = ( start_date.strftime("%Y-%m-%d") if isinstance(start_date, datetime.date) else start_date ) query_input["start_date"] = start_date_str result = StreamsByStore(query_input).execute() streams_timeseries = [ dict(row._mapping) if hasattr(row, "_mapping") else dict(row) for row in result ] streams_by_store = streams_utils.get_streams_by_store( streams_timeseries, start_date, end_date ) if streams_by_store: response_body["stores"] = streams_by_store schema = StreamsByStoreSchema() return oto.response.Response(schema_dump(schema, response_body))