import logging from statistics import mean from typing import List, Optional from fastapi import APIRouter, Body from service.api import schemas from service.tasks.api_service_handler.facebook.service_fb_db_handler import ( DbFbServiceHandler, ) from service.tasks.api_service_handler.facebook.service_fb_external import ( FbServiceHandlerExternal, ) from service.tasks.api_service_handler.facebook.service_fb_internal import ( FbServiceHandler, ) from service.tasks.api_service_handler.utils import field_mapper from service.tasks.audience import generate_fb_audience, initiate_fb_audience_sharing from service.tasks.data_model import list_caw_schemas from service.tasks.precondition.checks.deletion import FacebookAudienceByIdDeleteCheck from service.tasks.precondition.engine import Precondition from service.utils.appsync_communication.message_types import ( ServiceInfoError, ServiceUserTokenError, ) from service.utils.aws_connectors import run_query from service.utils.data_model_utils import get_rendered_sql_template logger = logging.getLogger(__name__) router = APIRouter(tags=["Facebook"]) @router.post("/get_fb_status", response_model=schemas.FbStatus) def get_fb_status( management_schema: str = Body(..., embed=True), ): db_io = DbFbServiceHandler() df = db_io.list_adaccounts(schema=management_schema) if len(df) > 0: result_response = {"status": "connected"} else: result_response = {"status": "not_connected"} # returns FbAccount appsync type return result_response @router.post("/get_fb_sdk_status", response_model=schemas.FbUser) @field_mapper def get_fb_sdk_status( management_schema: str = Body(..., embed=True), ): db_io = DbFbServiceHandler() df = db_io.list_users(schema=management_schema) if len(df) > 0: """Getting currently granted scopes""" try: fb_ext = FbServiceHandlerExternal(schema=management_schema) granted_permissions = fb_ext.get_granted_permissions() permission_list = [] for permission in granted_permissions["data"]["data"]: if permission["status"] == "granted": permission_list.append(permission["permission"]) db_io.update_fb_user_permissions( schema=management_schema, user_id=df["user_id"].values[0], permissions=permission_list, ) # returning FbUser appsync type result_response = { "status": "connected", "id": df["user_id"].values[0], "granted_scopes": permission_list, } except ServiceUserTokenError as e: db_io.delete_fb_user(schema=management_schema) raise ServiceInfoError(e.args[0]) except ServiceInfoError as e: raise ServiceInfoError(e.args[0]) else: result_response = { "status": "not_connected", "id": None, "granted_scopes": [], } return result_response @router.post("/list_fb_audiences", response_model=List[schemas.AudienceExtended]) @field_mapper def list_fb_audiences( user_id: str = Body(...), management_schema: str = Body(...), collectionIds: Optional[List[str]] = Body(None), adAccountIds: Optional[List[str]] = Body(None), ): """ New behavior: endpoint returns all the Alliances across all the workspaces and Alliances :param user_id: :param management_schema: :param collectionIds: :param adAccountIds: :param kw: :return: """ result = [] db_io = DbFbServiceHandler() if collectionIds is not None: pass else: for caw_schema, caw_type in list_caw_schemas( management_schema, with_caw_type=True ): df = db_io.list_audiences_and_adaccounts( schema=caw_schema, adaccount_id=adAccountIds, user_id=user_id, default_schema=management_schema, ) for rec in df.to_dict("records"): rec["source"] = {"cawType": caw_type, "cawSchema": caw_schema} result.append(rec) # Returning [AudienceExtended] type for appsync return result @router.post("/read_fb_audience", response_model=schemas.AudienceExtended) @field_mapper def read_fb_audience( audienceId: str = Body(...), user_id: str = Body(...), management_schema: str = Body(...), workspace_schema: Optional[str] = Body(None), alliance_schema: Optional[str] = Body(None), ): db_io = DbFbServiceHandler() for response in db_io.list_audiences_and_adaccounts( schema=alliance_schema or workspace_schema, internal_audience_id=audienceId, user_id=user_id, default_schema=management_schema, ).to_dict("records"): # Returning AudienceExtended type for appsync return response @router.post("/list_fb_adaccounts", response_model=List[schemas.AdAccount]) @field_mapper def list_fb_adaccounts( user_id: str = Body(...), management_schema: str = Body(...), ): """Updating list of adaccounts obtained through SDK""" update_fb_sdk_adaccounts(user_id, management_schema) """ Listing all fb adaccounts stored in db""" db_io = DbFbServiceHandler() df = db_io.list_adaccounts(schema=management_schema) # Returning [AdAccount] type for appsync return df.to_dict("records") @router.post("/list_saved_fb_sdk_adaccounts", response_model=List[schemas.AdAccount]) @field_mapper def list_saved_fb_sdk_adaccounts( management_schema: str = Body(..., embed=True), ): db_io = DbFbServiceHandler() df = db_io.list_saved_from_sdk_adaccounts(schema=management_schema) # Returning [AdAccount] appsync type response = df.to_dict("records") return response @router.post("/list_fb_sdk_adaccounts", response_model=List[schemas.AdAccount]) @field_mapper def list_fb_sdk_adaccounts_ep( management_schema: str = Body(...), user_id: str = Body(...), ): return list_fb_sdk_adaccounts( management_schema, user_id, ) def list_fb_sdk_adaccounts( management_schema: str, user_id: str, ): """Getting list of connected through SDK adaccounts""" db_io = DbFbServiceHandler() try: fb_ext = FbServiceHandlerExternal(schema=management_schema) response = fb_ext.list_user_owned_adaccounts() adaccounts_list = [] for adaccount in response["data"].get("data", []): business = adaccount.get("business", {"id": None, "name": None}) adaccount_obj = { "id": adaccount["id"].replace("act_", ""), "business_id": business["id"], "business_name": business["name"], "name": adaccount["name"], "status": adaccount["account_status"], "description": "", "user_id": user_id, } adaccounts_list.append(adaccount_obj) db_io = DbFbServiceHandler() db_io.clear_fb_sdk_adaccount_table(schema=management_schema) db_io.populate_fb_sdk_adaccount_table( schema=management_schema, adaccounts_list=adaccounts_list, fb_user_id=fb_ext.fb_user_id, user_id=user_id, ) except ServiceUserTokenError as e: db_io.delete_fb_user(schema=management_schema) raise ServiceInfoError(e.args[0]) except RuntimeError: """RuntimeError is fired when we don't have access_token User have not conencted facebook account through SDK """ adaccounts_list = [] db_io.clear_fb_sdk_adaccount_table(schema=management_schema) # returning [AdAccount] appsync type return adaccounts_list @router.post("/update_fb_sdk_adaccounts", response_model=List[schemas.AdAccount]) @field_mapper def update_fb_sdk_adaccounts_ep( management_schema: str = Body(...), user_id: str = Body(...), ): return update_fb_sdk_adaccounts(user_id, management_schema) def update_fb_sdk_adaccounts( user_id: str, management_schema: str, ): adaccounts_list = list_fb_sdk_adaccounts(management_schema, user_id) db_io = DbFbServiceHandler() db_io.store_adaccount_from_sdk(schema=management_schema) db_io.mark_disconnected_adaccounts_from_sdk(schema=management_schema) # returning [AdAccount] appsync type return adaccounts_list @router.post("/generate_fb_audience", response_model=schemas.AudienceExtended) @field_mapper def generate_fb_audience_ep( user_id: str = Body(...), collectionId: int = Body(...), audienceName: str = Body(...), audienceDescription: str = Body(...), management_schema: str = Body(...), workspace_schema: Optional[str] = Body(None), alliance_schema: Optional[str] = Body(None), ): for audience in generate_fb_audience( user_id, collectionId, audienceName, audienceDescription, management_schema, workspace_schema, alliance_schema, ): # returning AudienceExtended appsync type return audience raise RuntimeError("Audience not generated") @router.post("/delete_fb_audience", response_model=schemas.Audience) def delete_audience( user_id: str = Body(...), audienceId: str = Body(...), management_schema: str = Body(...), acceptedConsequences: List = Body(...), workspace_schema: Optional[str] = Body(None), alliance_schema: Optional[str] = Body(None), ): precondition = Precondition( FacebookAudienceByIdDeleteCheck( alliance_schema or workspace_schema, audienceId, management_schema=management_schema, user_id=user_id, ), accepted_consequences=acceptedConsequences, ) if precondition.is_met: precondition.execute() else: precondition.raise_exception() # Returning Audience type for appsync return { "id": audienceId, "name": "", "description": "", "size": 0, "collectionId": None, } @router.post("/store_fb_adaccount", response_model=schemas.AdAccount) @field_mapper def store_fb_adaccount( management_schema: str = Body(...), adAccountId: str = Body(...), adAccountName: Optional[str] = Body(None), adAccountDescription: Optional[str] = Body(None), ): """Storing adaccount data obtained manually""" db_io = DbFbServiceHandler() # returning Adaccount type for appsync return db_io.store_adaccount( schema=management_schema, adaccount_id=adAccountId, name=adAccountName, description=adAccountDescription, ) @router.post("/delete_fb_adaccount", response_model=schemas.AdAccount) def delete_fb_adaccount_ep( adAccountId: str = Body(...), management_schema: str = Body(...), workspace_schema: Optional[str] = Body(None), alliance_schema: Optional[str] = Body(None), ): return delete_fb_adaccount(adAccountId, management_schema) def delete_fb_adaccount( adAccountId: str, management_schema: str, ): # ! TODO - need to check and implement remove sharing agreement with business manager if no adaccounts left? db_io = DbFbServiceHandler() """ Unsharing audience with the adaccount that is being deleted""" shared_audiences_df = db_io.list_shared_audiences( adaccount_id=adAccountId, management_schema=management_schema, external_id=True ) external_audience_id_list = list(shared_audiences_df["external_id"]) fb = FbServiceHandler() for external_audience_id in external_audience_id_list: # TODO! This can be optimized for querying multiple rows fb.revoke_audience_sharing_from_adaccount( external_audience_id=external_audience_id, partner_ad_account_id=adAccountId ) # Deleting rows from audience_shared_state_table db_io.delete_from_fb_audience_shared_state_table( adaccount_id=adAccountId, management_schema=management_schema, everything=True ) """ Deleting adaccount from db""" db_io.delete_adaccount(schema=management_schema, adaccount_id=adAccountId) # returning Adaccount type for appsync return { "id": adAccountId, "name": None, "description": None, "status": None, "businessId": None, "businessName": None, } @router.post("/create_and_share_fb_audience", response_model=schemas.AudienceExtended) @field_mapper def create_and_share_fb_audience( user_id: str = Body(...), collectionId: int = Body(...), audienceName: str = Body(...), audienceDescription: str = Body(...), adAccountId: str = Body(...), management_schema: str = Body(...), workspace_schema: Optional[str] = Body(None), alliance_schema: Optional[str] = Body(None), ): response = generate_fb_audience( user_id, collectionId, audienceName, audienceDescription, management_schema, workspace_schema, alliance_schema, ) final_response = response[0] internal_audience_id = final_response["id"] # returns AudienceExtended appsync type return initiate_fb_audience_sharing( adAccountId, internal_audience_id, user_id=user_id, management_schema=management_schema, final_response=final_response, external_audience_id=response[1], workspace_schema=workspace_schema, alliance_schema=alliance_schema, ) @router.post("/list_shared_audiences", response_model=List[schemas.Audience]) @field_mapper def list_shared_audiences( adAccountId: str = Body(...), management_schema: str = Body(...), ): db_io = DbFbServiceHandler() df = db_io.list_shared_audiences(adAccountId, management_schema) # returns [Audience] appsync type return df.to_dict("records") @router.post("/generate_lookalike_fb_audience", response_model=schemas.AudienceExtended) @field_mapper def generate_lookalike_fb_audience( audienceName: str = Body(...), audienceDescription: str = Body(...), sourceAudienceId: str = Body(...), audienceRatio: float = Body(...), audienceCountryCodes: List[str] = Body(...), user_id: str = Body(...), management_schema: str = Body(...), workspace_schema: Optional[str] = Body(None), alliance_schema: Optional[str] = Body(None), ): db_io = DbFbServiceHandler() source_audience_data = db_io.get_audience_info( schema=alliance_schema or workspace_schema, internal_audience_id=sourceAudienceId, ) """Getting full FbCountry (appsync) objects""" countries = db_io.get_countries() countries = countries[countries["id"].isin(audienceCountryCodes)] countries = countries.to_dict("records") """Getting full FbLookalikeRatio (appsync) object""" ratios = db_io.get_lookalike_ratios() ratio = ratios[ratios["id"].isin([float(audienceRatio)])] ratio = ratio.to_dict("records")[0] """Creating lookalike audience (facebook api)""" fb = FbServiceHandler() audience_name = f"Fansifter - {audienceName}" response = fb.create_lookalike_audience( audience_name=audience_name, source_audience_id=source_audience_data["external_id"], audience_ratio=audienceRatio, audience_country_codes=audienceCountryCodes, audience_description=audienceDescription, ) # workaround made after introduction of lower and upper bounds in API v11 response["data"]["approximate_count"] = int( mean( [ response["data"]["approximate_count_lower_bound"], response["data"]["approximate_count_upper_bound"], ] ) ) # we don't have actual_count when it comes to lookalikes response["data"]["actual_count"] = response["data"]["approximate_count"] external_lookalike_audience_id = response["data"]["id"] """Populating corresponding DB tables with audience data""" db_io.populate_commons_fb_audience_meta_table( schema=alliance_schema or workspace_schema, external_audience_id=external_lookalike_audience_id, collection_id=source_audience_data["collection_id"], ) internal_lookalike_audience_id = db_io.populate_fb_audience_table( schema=alliance_schema or workspace_schema, response_obj=response["data"], collection_id=source_audience_data["collection_id"], parent_id=sourceAudienceId, countries=countries, ratio=ratio, user_id=user_id, ) # returning AudienceExtended type for appsync df = db_io.list_audiences_and_adaccounts( schema=alliance_schema or workspace_schema, internal_audience_id=internal_lookalike_audience_id, user_id=user_id, default_schema=management_schema, ) response = df.to_dict("records") try: response = response[0] except IndexError: response = None return response @router.post("/get_fb_country_suggestions", response_model=List[schemas.FbCountry]) def get_fb_country_suggestions(): db_io = DbFbServiceHandler() df = db_io.get_countries() # Returns [FbCountry] type for appsync return df.to_dict("records") @router.post( "/get_fb_lookalike_ratio_suggestions", response_model=List[schemas.FbLookalikeRatio] ) def get_fb_lookalike_ratio_suggestions(): db_io = DbFbServiceHandler() df = db_io.get_lookalike_ratios() # Returns [FbLookalikeRatio] type for appsync return df.to_dict("records") @router.post("/list_fb_sdk_required_scopes", response_model=List[str]) def list_fb_sdk_required_scopes(): """Returning a list of scopes that we need to request from FB user via SDk""" scopes = ["email", "ads_read"] return scopes @router.post("/initiate_fb_audience_sharing", response_model=schemas.AudienceExtended) @field_mapper def initiate_fb_audience_sharing_ep( user_id: str = Body(...), adAccountId: str = Body(...), audienceId: str = Body(...), management_schema: str = Body(...), workspace_schema: Optional[str] = Body(None), alliance_schema: Optional[str] = Body(None), ): """Initiating sharing process and actually sharing the audience if the sharing agreement with addaccount was previously established Adapted to be used in create_and_share_audience""" # returns AudienceExtended appsync type return initiate_fb_audience_sharing( adAccountId, audienceId, user_id=user_id, management_schema=management_schema, final_response=None, external_audience_id=None, workspace_schema=workspace_schema, alliance_schema=alliance_schema, ) @router.post("/update_fb_sdk_status", response_model=schemas.FbUser) @field_mapper def update_fb_sdk_status( user_id: str = Body(...), response: Optional[schemas.FbAuthInput] = Body(None), status: str = Body(...), management_schema: str = Body(...), ): db_io = DbFbServiceHandler() result = {} if status == "connected": if response is None: # todo split endpoint into two raise RuntimeError("response - field missing") fb = FbServiceHandler() long_lived_token_response = fb.get_long_lived_token( user_short_lived_token=response.accessToken ) fb_response = { **response.dict(), "accessToken": long_lived_token_response["data"]["access_token"], "expiresIn": long_lived_token_response["data"]["expires_in"], "token_type": long_lived_token_response["data"]["token_type"], } result_df = db_io.populate_fb_user_table( schema=management_schema, response=fb_response, status=status ) # Returing FbUser appsync type result = { "status": "connected", "id": result_df["user_id"].values[0], "granted_scopes": result_df["granted_permissions"].values[0].split(","), } elif status == "updated": """Not functional yet""" result = {"status": "updated", "id": None, "granted_scopes": []} elif status == "disconnected": """Request to delete the user data and disconnect fb_user_id from Fansifter""" fb_ext = FbServiceHandlerExternal(schema=management_schema) fb_ext.revoke_all_user_access() # Marking user and adaccounts for deletion upfront. # To avoid showing them in other queries db_io.mark_fb_user_for_deletion(schema=management_schema) db_io.mark_adaccounts_for_deletion( schema=management_schema, fb_user_id=fb_ext.fb_user_id ) db_response = db_io.get_adaccounts_by_user_id( schema=management_schema, user_id=fb_ext.fb_user_id ) try: for adaccount_id in list(db_response["id"]): delete_fb_adaccount( adAccountId=adaccount_id, management_schema=management_schema ) except (TypeError, KeyError): pass db_io.delete_fb_user(schema=management_schema) update_fb_sdk_adaccounts(user_id, management_schema) # Returing FbUser appsync type return {"status": "disconnected", "id": None, "granted_scopes": []} return result @router.post("/revoke_fb_adaccount_sharing", response_model=schemas.AudienceExtended) @field_mapper def revoke_fb_audience_sharing_from_adaccount( user_id: str = Body(...), audienceId: str = Body(...), adAccountId: str = Body(...), management_schema: str = Body(...), workspace_schema: Optional[str] = Body(None), alliance_schema: Optional[str] = Body(None), ): """Removing audience sharing with specific adaccount""" adaccount_id = str(adAccountId) db_io = DbFbServiceHandler() external_audience_id, account_id = db_io.get_audience_id( schema=alliance_schema or workspace_schema, internal_id=audienceId ) """FB API call to revoke sharing""" fb = FbServiceHandler(account_id) fb.revoke_audience_sharing_from_adaccount( external_audience_id=external_audience_id, partner_ad_account_id=adaccount_id ) """DB queries to remove data from table""" db_io.delete_from_fb_audience_shared_state_table( schema=alliance_schema or workspace_schema, internal_audience_id=audienceId, adaccount_id=adaccount_id, ) """Getting updated AudienceExtended type from db""" df = db_io.list_audiences_and_adaccounts( schema=alliance_schema or workspace_schema, internal_audience_id=audienceId, user_id=user_id, default_schema=management_schema, ) for final_response in df.to_dict("records"): # Returning AudienceExtended type for appsync return final_response return None @router.post("/update_fb_schema", response_model=schemas.FbSchemaUpdateStatus) def update_fb_schema( management_schema: str = Body(..., embed=True), ): params = {"schema_name": management_schema} db_io = DbFbServiceHandler() audiences_to_delete_df = db_io.list_audiences(schema=management_schema) if audiences_to_delete_df is not None: if ( len(audiences_to_delete_df) > 0 ): # Ensuring that there are audiences to delete audiences_data = [ {"internal_id": x[0], "external_id": x[1], "parent_id": x[2]} for x in audiences_to_delete_df[ ["id", "external_id", "parent_id"] ].itertuples(index=False) ] aundience_id_list = [ {"external_id": str(x["external_id"]), "parent_id": x["parent_id"]} for x in audiences_data ] fb = FbServiceHandler() fb.delete_multiple_audiences(audience_id_list=aundience_id_list) db_io.delete_multiple_audiences( schema=management_schema, audience_list=audiences_data ) """ Revoking access from fb_user""" try: fb_ext = FbServiceHandlerExternal(schema=management_schema) fb_ext.revoke_all_user_access() except Exception as e: logger.exception(f"FB revoke all user access error: {e}") """ Drop fb API related tables if exist + indexes""" template_sql = "drop_fb_tables.sql" sql = get_rendered_sql_template(params, template_sql) run_query(query_sql=sql) """ Generate fb base tables """ template_sql = "fb_integration/create_fb_base_tables.sql" sql = get_rendered_sql_template(params, template_sql) run_query(query_sql=sql) """ Generate fb company tables""" template_sql = "fb_integration/create_fb_company_tables.sql" sql = get_rendered_sql_template(params, template_sql) run_query(query_sql=sql) # returns FbSchemaUpdateStatus appsync type return {"status": "fb_tables_updated"}