import io import logging from collections.abc import Iterator from dataclasses import dataclass import pandas as pd from app.adapters.aws.kms import KMSClient from app.adapters.google import GoogleClient from app.adapters.ows_dmp import AudienceShare, AudienceSharePlatform, OwsDmpClient from app.exceptions import HandlerError from app.types import ProcessedFans from app.utils import smart_iter from .base import AudienceShareStrategy logger = logging.getLogger(__name__) @dataclass(kw_only=True) class GoogleAudienceShareStrategy(AudienceShareStrategy): ows_dmp_client: OwsDmpClient kms_client: KMSClient google_client: GoogleClient platform = AudienceSharePlatform.GOOGLE upload_chunksize = 10000 batch_size = 6 def share( self, audience_share: AudienceShare, file_buffer: io.BufferedRandom ) -> Iterator[ProcessedFans]: """Share Google audience and yield ProcessedFans.""" logger.info( "Fetching Google Audience Sharing information.", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, }, ) try: share_data = self.ows_dmp_client.get_google_audience_share( audience_share.id ) except Exception as exc: raise HandlerError( f"Failed to get Google Audience Sharing data `{audience_share.id}`.", code="GET_GOOGLE_AUDIENCE_SHARE_ERROR", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, }, ) from exc logger.info( "Decrypting Google Audience user access token.", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, "ad_account_external_id": share_data.ad_account_external_id, }, ) try: user_access_token = self.kms_client.decrypt( share_data.user_access_token, context={"field": "token"} ) except Exception as exc: raise HandlerError( "Failed to decrypt Google Audience User access token.", code="DECRYPT_USER_ACCESS_TOKEN_ERROR", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, "ad_account_external_id": share_data.ad_account_external_id, }, ) from exc # Upload Audience file logger.info( "Reading Audience file.", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, "ad_account_external_id": share_data.ad_account_external_id, }, ) df = pd.read_csv(file_buffer, chunksize=self.upload_chunksize) total_batches = share_data.fans_count // self.upload_chunksize + 1 logger.info( "Creating offline user data Job", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, "ad_account_external_id": share_data.ad_account_external_id, }, ) offline_job = self.google_client.create_offline_user_data_job( access_token=user_access_token, customer_id=share_data.ad_account_external_id, user_list_resource=share_data.user_list_resource, login_customer_id=share_data.parent_ad_account_external_id, ) for batch_seq, is_last, data in smart_iter(df, start=1): processed_fans = data.shape[0] self.google_client.add_to_offline_user_data_job( access_token=user_access_token, offline_job_resource=offline_job.resource_name, payload=data[["FAN_ID", "FAN_PHONE"]].fillna("").values.tolist(), login_customer_id=share_data.parent_ad_account_external_id, ) logger.info( "Uploaded Google Audience data - %s of %s.", batch_seq, total_batches, extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, "ad_account_external_id": share_data.ad_account_external_id, "upload_batch_seq": batch_seq, "upload_last_batch_flag": is_last, "processed_fans": processed_fans, }, ) yield ProcessedFans(processed_fans=processed_fans, is_last=is_last) logger.info( "Starting offline job", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, "ad_account_external_id": share_data.ad_account_external_id, }, ) self.google_client.start_offline_user_data_job( access_token=user_access_token, offline_job_resource=offline_job.resource_name, login_customer_id=share_data.parent_ad_account_external_id, )