from datetime import timedelta, date from typing import List, Dict, Optional from data_health.data_health_repository import DataHealthRepository from data_health.schemas import ( PerformanceMetricsDataHealthSchema, DataHealthResponseUnit, PerformanceMetricsDataHealthResponse, ReportingDataHeathQuery, ReportingDataHealthResponse, ) from services.permissions.exceptions import DataHealthNotAvailable from models.marketing_accounts import MarketingAccount from projects.repositories.projects_repository import ProjectsRepository from utils.list_utils import get_nested_attr from external_api.base.clients.client_factory import ApiClientFactory from external_api.base.clients.delphi_client import ( DelphiClient, PerformanceMetricsDataHealthQuery, PerformanceMetricsDataHealthDelphiResponse, ReportingDataHeathResponse, ) statuses_mapping = { None: 1, "missing": 1, "pending": 1, "complete": 2, "complete_minimal": 2, "complete_volatile": 2, "partial": 3, } RESPONSE_STATUS_MISSING = 1 RESPONSE_STATUS_COMPLETE = 2 RESPONSE_STATUS_PARTIAL = 3 reverse_statuses_mapping = { RESPONSE_STATUS_MISSING: "missing", RESPONSE_STATUS_COMPLETE: "complete", RESPONSE_STATUS_PARTIAL: "partial" } class DataHealthService: data_health_repository: DataHealthRepository = DataHealthRepository() delphi_client: DelphiClient = ApiClientFactory.delphi_client() projects_repository: ProjectsRepository = ProjectsRepository() @staticmethod def get_external_id(account: MarketingAccount): return account.source + "_" + account.external_id async def get_performance_metrics_data_health(self, params: PerformanceMetricsDataHealthSchema): marketing_accounts = self.data_health_repository.get_marketing_accounts_by_label_ids(params.labelIds) ma_external_ids = [self.get_external_id(ma) for ma in marketing_accounts] delphi_query = PerformanceMetricsDataHealthQuery( start_date=params.startDate, end_date=params.endDate, ads_account=ma_external_ids ) data = await self.delphi_client.get_performance_metrics_data_health(delphi_query) # todo get Delphi to always send us the dates data.min_date = data.min_date or params.startDate data.max_date = data.max_date or params.endDate return self.__map_data_health_response(data, marketing_accounts) def __map_data_health_response( self, data: PerformanceMetricsDataHealthDelphiResponse, marketing_accounts: List[MarketingAccount] ) -> PerformanceMetricsDataHealthResponse: linkfire_days = get_nested_attr(data, "breakdowns.dsp.linkfire.days", None) facebook_data = get_nested_attr(data, "breakdowns.dsp_parent_rep_owner_ads_account.facebook", None) google_data = get_nested_attr(data, "breakdowns.dsp_parent_rep_owner_ads_account.google", None) result = PerformanceMetricsDataHealthResponse() result.linkfire = self.__map_dsp_range(linkfire_days, data.min_date, data.max_date) for name, dsp_data in {"facebook": facebook_data, "google": google_data}.items(): if dsp_data is None: # todo get Delphi to always send us complete data setattr( result, name, [ DataHealthResponseUnit( startDate=data.min_date, endDate=data.max_date, status=reverse_statuses_mapping.get(RESPONSE_STATUS_MISSING), adAccounts=set([self.get_external_id(a) for a in marketing_accounts if a.source == name]), ) ] ) else: setattr(result, name, self.__map_ads_accounts(dsp_data, data.min_date)) return result def __prepare_dsp_ad_accounts_data(self, dsp_data: Dict) -> List: result = [] for label_id in dsp_data.keys(): for ad_account_id in dsp_data[label_id].keys(): result.append({"key": ad_account_id, "val": dsp_data[label_id][ad_account_id]["days"]}) return result def __map_ads_accounts(self, statuses_by_day: Dict, min_date: date) -> List[DataHealthResponseUnit]: statuses_by_day = self.__prepare_dsp_ad_accounts_data(statuses_by_day) if not statuses_by_day: return [] result = [DataHealthResponseUnit(startDate=min_date, endDate=min_date, adAccounts=set(), status=None)] def is_ad_account_missing_data(ad_account_data): return ad_account_data in (RESPONSE_STATUS_MISSING, RESPONSE_STATUS_PARTIAL) for i in range(len(statuses_by_day[0]["val"])): day_status_binary = 0 ad_accounts_set = set() for item in statuses_by_day: ad_account_status = statuses_mapping[item["val"][i]] if is_ad_account_missing_data(ad_account_status): ad_accounts_set.add(item["key"]) day_status_binary |= ad_account_status day_status = reverse_statuses_mapping[day_status_binary] result[-1].status = result[-1].status or day_status if day_status == result[-1].status: result[-1].endDate = min_date result[-1].adAccounts.update(ad_accounts_set) else: result.append( DataHealthResponseUnit( startDate=min_date, endDate=min_date, status=day_status, adAccounts=ad_accounts_set ) ) min_date += timedelta(days=1) return result async def get_project_data_heath_for_reporting_tab(self, project_id: int, params: ReportingDataHeathQuery): product_families = self.projects_repository.get_product_families(project_id) tracks = [] if product_families is not None: for product_family in product_families: if product_family.tracks is not None: tracks += product_family.tracks playlists = self.projects_repository.get_target_playlists if len(tracks) == 0 or playlists is None: raise DataHealthNotAvailable() return await self.__get_reporting_data_heath(params) async def __get_reporting_data_heath(self, params: ReportingDataHeathQuery): data = await self.delphi_client.get_reporting_tab_data_health( start_date=params.startDate, end_date=params.endDate ) return self.__map_reporting_data_health(data, params.startDate, params.endDate) def __map_reporting_data_health( self, response: ReportingDataHeathResponse, start_date: date, end_date: date, ) -> ReportingDataHealthResponse: min_date = response.min_date or start_date max_date = response.max_date or end_date result = ReportingDataHealthResponse( amazon=self.__map_dsp_range( get_nested_attr(response, "breakdowns.dsp_segment.amazon.streams.days", None), min_date, max_date ), apple=self.__map_dsp_range( get_nested_attr(response, "breakdowns.dsp_segment.apple.streams.days", None), min_date, max_date ), spotify=self.__map_dsp_range( get_nested_attr(response, "breakdowns.dsp_segment.spotify.streams.days", None), min_date, max_date ), ) return result def __map_dsp_range(self, statuses_by_day: Optional[Dict], min_date: date, max_date: date): if not statuses_by_day: return [DataHealthResponseUnit( startDate=min_date, endDate=max_date, status=reverse_statuses_mapping.get(RESPONSE_STATUS_MISSING), )] status = self.__get_day_status(statuses_by_day.pop(0)) result = [DataHealthResponseUnit(startDate=min_date, endDate=min_date, status=status)] for item in statuses_by_day: item = self.__get_day_status(item) min_date += timedelta(days=1) if item == result[-1].status: result[-1].endDate = min_date else: result.append(DataHealthResponseUnit(startDate=min_date, endDate=min_date, status=item)) return result def __get_day_status(self, status): return reverse_statuses_mapping[statuses_mapping[status]]