import logging import tempfile from collections.abc import Iterator from dataclasses import dataclass from app.adapters.aws.s3 import S3Client from app.adapters.ows_dmp import ( AudienceShare, AudienceShareStatus, OwsDmpClient, ) from app.exceptions import HandlerError from app.strategies import AudienceShareStrategy from app.types import ProcessedFans logger = logging.getLogger(__name__) @dataclass class ShareAudienceRequest: bucket: str key: str @dataclass(kw_only=True) class ShareAudienceHandler: s3_client: S3Client ows_dmp_client: OwsDmpClient strategies: list[AudienceShareStrategy] def handle(self, request: ShareAudienceRequest) -> None: # Get audience share id and filename by S3 Key *_, audience_share_id, _ = request.key.split("/") logger.info( "Fetching Audience Share Information.", extra={"audience_share_id": audience_share_id}, ) # Set share status 'PROCESSING' audience_share = self._update_status( audience_share_id, status=AudienceShareStatus.PROCESSING, ) total_processed_fans = 0 for processed_fans, is_last in self._process( request, audience_share=audience_share ): total_processed_fans += processed_fans if is_last: status = AudienceShareStatus.COMPLETED else: status = AudienceShareStatus.PROCESSING # Update share status self._update_status( audience_share, status=status, processed_fans=total_processed_fans, ) return None @property def _strategy_by_platform(self) -> dict[str, AudienceShareStrategy]: return {strategy.platform: strategy for strategy in self.strategies} def _process( self, request: ShareAudienceRequest, *, audience_share: AudienceShare ) -> Iterator[ProcessedFans]: """Process audience and yield ProcessedFans.""" logger.info( "Processing Audience...", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, }, ) try: yield from self._share_audience(request, audience_share=audience_share) except Exception as exc: self._update_status(audience_share, status=AudienceShareStatus.FAILED) if isinstance(exc, HandlerError): raise exc raise HandlerError( "Failed to process Audience Sharing.", code="PROCESS_AUDIENCE_ERROR", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, }, ) from exc def _share_audience( self, request: ShareAudienceRequest, audience_share: AudienceShare ) -> Iterator[ProcessedFans]: """Share audience by supported platform and yield ProcessedFans.""" try: strategy = self._strategy_by_platform[audience_share.platform] except KeyError as exc: raise HandlerError( f"Unsupported Audience Sharing Platform `{audience_share.platform}`.", code="UNSUPPORTED_AUDIENCE_PLATFORM", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, }, ) from exc with tempfile.TemporaryFile() as file_obj: logger.info( "Downloading Audience file", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, "s3_bucket": request.bucket, "s3_key": request.key, }, ) try: self.s3_client.download_fileobj( bucket=request.bucket, key=request.key, file_obj=file_obj, ) file_obj.seek(0) except Exception as exc: raise HandlerError( "Failed to download Audience file.", code="DOWNLOAD_AUDIENCE_FILE_ERROR", extra={ "audience_share_id": audience_share.id, "audience_platform": audience_share.platform, "s3_bucket": request.bucket, "s3_key": request.key, }, ) from exc yield from strategy.share(audience_share, file_obj) def _update_status( self, audience_share: AudienceShare | str, /, status: AudienceShareStatus, processed_fans: int | None = None, ) -> AudienceShare: """Update share status and return updated share.""" if isinstance(audience_share, AudienceShare): audience_share_id = audience_share.id audience_platform = audience_share.platform else: audience_share_id = audience_share audience_platform = None try: return self.ows_dmp_client.update_audience_share( audience_share_id, status=status, processed_fans=processed_fans, ) except Exception as exc: raise HandlerError( f"Failed to update Audience Share `{audience_share_id}` " f"status to `{status}`.", code="UPDATE_AUDIENCE_STATUS_ERROR", extra={ "audience_share_id": audience_share_id, "audience_platform": audience_platform, "status": status, }, ) from exc