import datetime import logging from dataclasses import dataclass from audience_common import datetimeutc from audience_common.auth.account import Account from campaigns.auth.requests import AuthRequest from campaigns.auth.services import AuthService from campaigns.config import settings from campaigns.connectors.aws.kms import BaseKMS from campaigns.connectors.db import Database, session_context from campaigns.connectors.facebook import models as api_models from campaigns.connectors.facebook.base import FacebookClient from campaigns.connectors.facebook.enums import ErrorCode from campaigns.connectors.facebook.exceptions import FacebookClientError from campaigns.connectors.facebook.utils import get_page_picture_content_in_bulk from campaigns.core.handler import Handler from campaigns.core.types import EncryptedToken, PlainToken from campaigns.meta import constants from campaigns.meta.dtos import FacebookConnection, FacebookUser from campaigns.meta.enums import FacebookConnectionStatus from campaigns.meta.models import FacebookAssociation, FacebookPage, FacebookPageLabel from campaigns.meta.repositories import ( FacebookAssociationRepository, FacebookPageLabelRepository, ) from campaigns.meta.services import FacebookPageService logger = logging.getLogger(__name__) @dataclass class ConnectFacebookUserRequest(AuthRequest): user_id: str access_token: PlainToken @dataclass class ConnectFacebookUserHandler( Handler[ConnectFacebookUserRequest, FacebookConnection], ): db: Database kms: BaseKMS auth_service: AuthService facebook_client: FacebookClient facebook_association_repository: FacebookAssociationRepository facebook_page_label_repository: FacebookPageLabelRepository facebook_page_service: FacebookPageService @session_context(in_transaction=True) async def handle(self, request: ConnectFacebookUserRequest) -> FacebookConnection: account_access = await self.auth_service.authorize_account(request.profile_id) association, created = await self._get_or_create_association( identity_id=request.identity_id ) api_token = await self.facebook_client.get_oauth_access_token( exchange_token=request.access_token, ) api_user = await self.facebook_client.get_user( user_id=request.user_id, user_access_token=api_token.access_token, ) association.user_id = request.user_id association.user_name = api_user.name association.token = await self.kms.encrypt( api_token.access_token, context={"field": "token"} ) association.is_connected = True association.expires_at = datetimeutc.now() + datetime.timedelta( seconds=api_token.expires_in or constants.DEFAULT_EXPIRES_IN ) # Get Facebook page from Graph API api_pages = await self._get_api_pages( user_id=request.user_id, user_access_token=api_token.access_token, ) # Return connection status if there are no available pages for current user if not api_pages and created: return await self._save_connected(association, update_user=api_user) api_encrypted_tokens = await self.kms.encrypt_many( [api_page.access_token for api_page in api_pages], context={"field": "page_token"}, ) api_decrypted_tokens = { encrypted_token: plain_token for plain_token, encrypted_token in api_encrypted_tokens.items() } await self._create_or_update_pages( association, api_pages=api_pages, api_encrypted_tokens=api_encrypted_tokens ) if association.pages: # we share all Facebook pages with the User's Vendor/Subaccount # automatically only if the user has direct access to a single one if len(account_access.accounts) == 1: await self._create_page_labels( association, account=account_access.accounts[0] ) if settings.facebook_assign_agency_access_to_pages_enabled: await self._assign_agency_access_for_pages( association, api_decrypted_tokens=api_decrypted_tokens ) return await self._save_connected(association, update_user=api_user) async def _get_or_create_association( self, identity_id: str ) -> tuple[FacebookAssociation, bool]: association = ( await self.facebook_association_repository.get_by_identity_id_or_none( identity_id=identity_id ) ) if not association: association = FacebookAssociation() association.orchard_identity_id = identity_id return association, True return association, False async def _create_or_update_pages( self, association: FacebookAssociation, *, api_pages: list[api_models.Page], api_encrypted_tokens: dict[PlainToken, EncryptedToken], ) -> None: api_page_by_id = {api_page.id: api_page for api_page in api_pages} page_by_id = {page.external_id: page for page in association.pages} existing_ids = set(page_by_id.keys()) actual_ids = set(api_page_by_id.keys()) created_ids = actual_ids.difference(existing_ids) updated_ids = existing_ids.intersection(actual_ids) deleted_ids = existing_ids.difference(actual_ids) picture_contents = await get_page_picture_content_in_bulk(api_pages) for page_id in created_ids: api_page = api_page_by_id[page_id] page = FacebookPage.from_api_obj(api_page) page.token = api_encrypted_tokens[api_page.access_token] if picture_content := picture_contents.get(page_id): await self.facebook_page_service.save_page_picture( page, content=picture_content ) association.pages.append(page) for page_id in updated_ids: api_page = api_page_by_id[page_id] # Update existing database page by existing Facebook API page page = page_by_id[page_id] page.token = api_encrypted_tokens[api_page.access_token] page.is_manageable = api_page.is_manageable page.is_valid = True page.is_agency_provider = False # # Update the existing S3 picture if picture_content := picture_contents.get(page_id): await self.facebook_page_service.save_page_picture( page, content=picture_content ) for page_id in deleted_ids: page = page_by_id[page_id] page.is_valid = False page.is_agency_provider = False # TODO: action required for campaigns which have invalid page ids async def _create_page_labels( self, association: FacebookAssociation, *, account: Account ) -> None: existing_page_labels = ( await self.facebook_page_label_repository.find_by_page_ids_and_account( page_ids=[page.id for page in association.pages], account=account ) ) existing_page_label_by_page_id = { page_label.page_id: page_label for page_label in existing_page_labels } for page in association.pages: if page.id not in existing_page_label_by_page_id: page_label = FacebookPageLabel() page_label.account = account page.labels.append(page_label) async def _assign_agency_access_for_pages( self, association: FacebookAssociation, api_decrypted_tokens: dict[EncryptedToken, PlainToken], ) -> None: manageable_pages = [ page for page in association.pages if page.is_manageable and page.is_valid ] for page in manageable_pages: page_access_token = api_decrypted_tokens[page.token] try: await self.facebook_client.assign_agency_to_business_manager( page_id=page.external_id, page_access_token=page_access_token, ) except FacebookClientError as exc: # DUPLICATE_ASSET_ASSIGNMENT code is a working scenario # when we try to assign agency access for a page that we # already have agency access (provided by other user) if exc.fb_code == ErrorCode.DUPLICATE_ASSET_ASSIGNMENT: continue logger.error( "Could not assign agency to business manager.", extra=exc.get_log_extra(), ) page.is_manageable = False continue try: await self.facebook_client.assign_agency_to_app_system_user( page.external_id, ) except FacebookClientError as exc: logger.error( "Could not assign agency to app system user.", extra=exc.get_log_extra(), ) page.is_manageable = False continue page.is_agency_provider = True async def _save_connected( self, association: FacebookAssociation, update_user: api_models.User ) -> FacebookConnection: self.facebook_association_repository.add(association) return FacebookConnection( user=FacebookUser.from_api_obj(update_user), status=FacebookConnectionStatus.CONNECTED, ) async def _get_api_pages( self, user_id: str, user_access_token: PlainToken ) -> list[api_models.Page]: result = [] async for account in self.facebook_client.get_user_pages( user_id, user_access_token=user_access_token ): result.append(account) return result