import datetime 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.facebook import FacebookClient from app.adapters.facebook.types import UploadUsersPayload from app.adapters.ows_dmp import AudienceShare, AudienceSharePlatform, OwsDmpClient from app.exceptions import HandlerError from app.types import ProcessedFans from .base import AudienceShareStrategy logger = logging.getLogger(__name__) @dataclass(kw_only=True) class MetaAudienceShareStrategy(AudienceShareStrategy): ows_dmp_client: OwsDmpClient kms_client: KMSClient facebook_client: FacebookClient platform = AudienceSharePlatform.META upload_chunksize = 10000 batch_size = 10 def share( self, audience_share: AudienceShare, file_buffer: io.BufferedRandom ) -> Iterator[ProcessedFans]: logger.info( "Fetching Meta Audience Sharing information.", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, }, ) try: share_data = self.ows_dmp_client.get_meta_audience_share(audience_share.id) except Exception as exc: raise HandlerError( f"Failed to get Meta Audience Sharing data `{audience_share.id}`.", code="GET_META_AUDIENCE_SHARE_ERROR", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, }, ) from exc logger.info( "Decrypting Meta Audience user access token.", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, "audience_external_id": share_data.audience_external_id, "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 Meta Audience User access token.", code="DECRYPT_USER_ACCESS_TOKEN_ERROR", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, "audience_external_id": share_data.audience_external_id, "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, "audience_external_id": share_data.audience_external_id, "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 session_id = int( datetime.datetime.timestamp(datetime.datetime.now(datetime.UTC)) * 10 ) def _upload_batch( data: pd.DataFrame, batch_seq: int, is_last: bool ) -> ProcessedFans: processed_fans = data.shape[0] payload = UploadUsersPayload( schema=["EMAIL", "PHONE"], data=data[["FAN_ID", "FAN_PHONE"]].fillna("").values.tolist(), ) logger.info( "Uploading Meta 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, }, ) self.facebook_client.upload_users( user_access_token, audience_external_id=share_data.audience_external_id, session={ "session_id": session_id, "batch_seq": batch_seq, "last_batch_flag": is_last, "estimated_num_total": share_data.fans_count, }, payload=payload, ) return ProcessedFans(processed_fans=processed_fans, is_last=is_last) # Process and upload data in parallel yield from self.process_in_parallel(df=df, process_fn=_upload_batch)