import copy import json from typing import List import numpy as np import pandas as pd from psycopg2.extensions import AsIs, register_adapter from sqlalchemy.exc import IntegrityError from service.tasks.data_model import list_caw_schemas from service.utils.aws_connectors import df_to_db, run_query class DbFbServiceHandler: """Separate class handling all the db communications for FbServiceHandler""" def __init__(self): """Setting up adapters for DB functions to understand Numpy's data types""" def adapt_numpy_float64(numpy_float64): return AsIs(numpy_float64) def adapt_numpy_int64(numpy_int64): return AsIs(numpy_int64) register_adapter(np.float64, adapt_numpy_float64) register_adapter(np.int64, adapt_numpy_int64) @staticmethod def populate_commons_fb_audience_meta_table( schema, external_audience_id, collection_id, public_id=None ): """Handling formatting and populating the fb_audience_meta table""" df = pd.DataFrame( { "audience_id": [external_audience_id], "schema_name": [schema], "collection_id": [collection_id], "public_id": [public_id], } ) df_to_db(df=df, schema="commons", table_name="fb_audience_meta") return True @staticmethod def populate_fb_audience_table( schema, response_obj, collection_id, countries=None, ratio=None, parent_id=None, user_id=None, ): """Handling formatting and populating the fb_audience table""" df = pd.DataFrame(response_obj) df.reset_index(inplace=True) df["operation_status_code"] = df[df["index"] == "code"]["operation_status"] df["operation_status_description"] = df[df["index"] == "description"][ "operation_status" ] df["user_id"] = user_id df.fillna(method="bfill", inplace=True) df = df[df["index"] == "code"] result_columns = [ "id", "account_id", "user_id", "approximate_count", "actual_count", "description", "name", "operation_status_code", "operation_status_description", ] df = df[result_columns] df["collection_id"] = collection_id now = pd.Timestamp.now() df["timestamp_created"] = now df["timestamp_last_updated"] = now """Lookalike Audiences always have parent_id, countries and ratio""" if parent_id is not None: df["parent_id"] = parent_id df["ratio"] = json.dumps(ratio) df["countries"] = json.dumps(countries) df.rename(columns={"id": "external_id"}, inplace=True) df = df.iloc[0] # Ensuring that only one row is provided params = df.to_dict() insert_fields = '("' + '", "'.join([key for key in params]) + '")' insert_values = "(" + ", ".join(["%(" + key + ")s" for key in params]) + ")" query = f""" INSERT INTO {schema}.fb_audience {insert_fields} VALUES {insert_values} RETURNING id;""" internal_audience_id = run_query(query, params)[0][0] return internal_audience_id @staticmethod def delete_audience(schema, internal_audience_id, external_audience_id): params = { "internal_id": internal_audience_id, "audience_id": external_audience_id, } query = f""" DELETE FROM {schema}.fb_audience WHERE id = %(internal_id)s; """ run_query(query, params) query = """ DELETE FROM commons.fb_audience_meta WHERE audience_id = %(audience_id)s; """ run_query(query, params) query = f""" DELETE FROM {schema}.fb_audience_shared_state WHERE audience_id = %(internal_id)s; """ run_query(query, params) def delete_multiple_audiences(self, schema, audience_list): for audience in audience_list: self.delete_audience( schema=schema, internal_audience_id=audience["internal_id"], external_audience_id=audience["external_id"], ) @staticmethod def list_audiences(schema) -> pd.DataFrame: query = f""" SELECT id, external_id, name, description, actual_count as "size", collection_id, parent_id FROM {schema}.fb_audience; """ df = run_query(query_sql=query, return_type="df") if df is None: df = pd.DataFrame( { "id": [], "external_id": [], "name": [], "description": [], "size": [], "collection_id": [], "parent_id": [], } ) return df def list_audiences_and_adaccounts( self, schema, user_id, default_schema, internal_audience_id=None, adaccount_id: List[str] = None, ) -> pd.DataFrame: """Getting list of all audiences we have in the schema Returning AudienceExtended appsync type""" if internal_audience_id is not None: where_statement = f"WHERE id = {int(internal_audience_id)}" elif adaccount_id is not None: adaccount_id_list = "', '".join([str(int(x)) for x in adaccount_id]) where_statement = ( f"WHERE id IN (SELECT DISTINCT audience_id " f"FROM {schema}.fb_audience_shared_state " f"WHERE adaccount_id IN ('{adaccount_id_list}'))" ) else: where_statement = "WHERE 1=1" query = f""" SELECT id, name, description, actual_count as "size", collection_id, parent_id, ratio, countries, TO_CHAR(timestamp_created, 'YYYY-MM-DD"T"HH24:MI:SS.US"Z"') as "timestamp_created", TO_CHAR(timestamp_last_updated, 'YYYY-MM-DD"T"HH24:MI:SS.US"Z"') as "timestamp_last_updated" FROM {schema}.fb_audience {where_statement} AND user_id = '{user_id}'; """ df = run_query(query_sql=query, return_type="df") """Returning empty df in case if query did not work or no results in db""" if df is None or len(df) < 1: df = pd.DataFrame( { "id": [], "name": [], "description": [], "size": [], "collection_id": [], "parent_id": [], "ratio": [], "countries": [], "timestamp_created": [], "timestamp_last_updated": [], } ) else: """Getting list of all shared adaccounts converting in to list of dicts per audience_id (parent_id)""" df2 = self.list_shared_adaccounts( schema=schema, default_schema=default_schema ) if len(df2) < 1: """In case if we don't have a single adaccount saved yet We need to force df2 to have at least one row in order for the further logic to work""" df2 = pd.DataFrame({"parent_id": [np.nan], "shared_with": [np.nan]}) else: df2 = ( df2.groupby("parent_id")[ "id", # AdAccountExtended type for appsync "adaccount_id", "name", "description", "shared_status", "status", "business_id", "business_name", ] .apply(pd.DataFrame.to_dict, "records") .reset_index() ) df2.rename(columns={0: "shared_with"}, inplace=True) df = df.merge(df2, how="left", left_on="id", right_on="parent_id") """Audience data already has parent_id which is referring to Lookalike audience -> custom audience relation""" df.rename(columns={"parent_id_x": "parent_id"}, inplace=True) """Filling null values (when audience is not shared) with empty lists for front end""" df["list_filler"] = np.empty((len(df), 0)).tolist() df["shared_with"] = ( df[["shared_with", "list_filler"]] .fillna(method="bfill", axis=1) .iloc[:, 0] ) """Filling null values (when no countries are available - custom Audience) with empty lists for front end""" df["countries"] = ( df[["countries", "list_filler"]] .fillna(method="bfill", axis=1) .iloc[:, 0] ) df.drop(columns=["parent_id_y", "list_filler"], inplace=True) return df @staticmethod def list_saved_from_sdk_adaccounts(schema): """Listing all adaccounts that came from currently saved fb_user_id from SDK""" query = f""" SELECT id, name, description, status, business_id, business_name FROM {schema}.fb_adaccount WHERE fb_user_id IS NOT NULL AND fb_user_id IN (SELECT DISTINCT user_id FROM {schema}.fb_user WHERE status != 'deleting') AND status != 907; """ df = run_query(query, return_type="df") return df @staticmethod def list_shared_adaccounts(schema, default_schema, internal_audience_id=None): """Getting list of adaccounts that have at least one audience shared with them OR that are shared with specific audience_id""" if internal_audience_id is not None: where_statement = f"AND audience_id = {int(internal_audience_id)}" else: where_statement = "" query = f""" SELECT fass.adaccount_id || '-' || fass.audience_id id, fass.adaccount_id adaccount_id, fa.name, fa.description, fa.status, fa.business_id, fa.business_name, fass.shared_status, fass.audience_id parent_id FROM {schema}.fb_audience_shared_state fass LEFT JOIN {default_schema}.fb_adaccount fa ON fass.adaccount_id = fa.id WHERE fa.status != 907 {where_statement}; """ df = run_query(query_sql=query, return_type="df") if df is None or len(df) < 1: empty_data = { "id": [], "adaccount_id": [], "name": [], "description": [], "status": [], "business_id": [], "business_name": [], "shared_status": [], "parent_id": [], } df = pd.DataFrame(empty_data) return df @staticmethod def populate_fb_sdk_adaccount_table( schema, user_id, adaccounts_list: List[dict], fb_user_id: str ) -> None: """Storing provided adacocunts in to temporary fb_sdk_adaccount table Manually adding created and updated timestamps and user_id """ if len(adaccounts_list) == 0: return adaccounts_copy_list = copy.deepcopy(adaccounts_list) insert_fields = ( '("' + '", "'.join([key for key in adaccounts_copy_list[0]]) + '", "timestamp_created", "timestamp_updated", "fb_user_id")' ) insert_values_list = [] params = {} now = pd.Timestamp.now() for index, adaccount in enumerate(adaccounts_copy_list): adaccount["timestamp_created"] = str(now) adaccount["timestamp_updated"] = str(now) adaccount["fb_user_id"] = fb_user_id insert_values_list.append( "(" + ", ".join(["%(" + key + str(index) + ")s" for key in adaccount]) + ")" ) for key in adaccount: final_key = key + str(index) params[final_key] = adaccount[key] insert_values = ", ".join(insert_values_list) query = f""" INSERT INTO {schema}.fb_sdk_adaccount {insert_fields} VALUES {insert_values} RETURNING *; """ run_query(query, params, return_type="df") @staticmethod def clear_fb_sdk_adaccount_table(schema): query = f""" DELETE FROM {schema}.fb_sdk_adaccount; """ run_query(query) @staticmethod def get_audience_info(schema, internal_audience_id): """Getting only the audience data without shared adaccounts""" params = {"internal_audience_id": internal_audience_id} query = f""" SELECT id, external_id, name, description, actual_count as "size", collection_id, parent_id, ratio, countries FROM {schema}.fb_audience WHERE id = %(internal_audience_id)s; """ df = run_query(query, params, return_type="df") response = df.iloc[0].to_dict() return response @staticmethod def get_audience_id(schema, internal_id): for aud in run_query( f"SELECT external_id id, account_id FROM {schema}.fb_audience WHERE id = %(audience_id)s;", {"audience_id": int(internal_id)}, ): return aud[0], aud[1] raise RuntimeError("Invalid id") @staticmethod def get_audience_collection_id(schema, internal_id): for coll in run_query( f"SELECT collection_id id FROM {schema}.fb_audience WHERE id = %(audience_id)s;", {"audience_id": int(internal_id)}, ): return coll[0] raise RuntimeError("Invalid id") @staticmethod def store_adaccount(schema, adaccount_id, name=None, description=None): if name is None: name = "No Name" if description is None: description = "No description" data = { "id": [adaccount_id], "name": [name], "description": [description], "status": 908, } df = pd.DataFrame(data) now = pd.Timestamp.now() df["timestamp_created"] = now df["timestamp_updated"] = now df["user_id"] = schema try: df_to_db(df=df, schema=schema, table_name="fb_adaccount") except IntegrityError: pass # returning Adaccount type for appsync response = { "id": adaccount_id, "name": name, "description": description, "status": None, "business_id": None, "business_name": None, } return response @staticmethod def store_adaccount_from_sdk( schema, adaccount_id=None, description=None ) -> pd.DataFrame: """Copies obtained data from the SDK (fb_sdk_adaccount) into the final adaccount table (fb_adaccount) On constraint updates the fields instead""" params = {"adaccount_id": adaccount_id, "description": description} update_description = "" where_section = "" if description is not None: update_description = ( "description = coalesce(a.description, %(description)s)," ) if adaccount_id is not None: where_section = "WHERE b.id = %(adaccount_id)s" query = f""" INSERT INTO {schema}.fb_adaccount AS a SELECT * FROM {schema}.fb_sdk_adaccount AS b {where_section} ON CONFLICT ON CONSTRAINT unique_fb_adaccount_id DO UPDATE SET business_id = excluded.business_id, business_name = excluded.business_name, fb_user_id = excluded.fb_user_id, name = excluded.name, {update_description} status = coalesce(excluded.status, a.status), timestamp_updated = excluded.timestamp_created, user_id = excluded.user_id RETURNING id, name, description, status, business_id, business_name; """ df = run_query(query, params, return_type="df") try: return df.to_dict("records")[0] except IndexError: """In case there are no stored SDK adaccounts""" return [] except AttributeError: # In case None data farame. return [] @staticmethod def mark_disconnected_adaccounts_from_sdk(schema): """Sets the status of disconnected adaccount to 909""" now = pd.Timestamp.now() query = f""" UPDATE {schema}.fb_adaccount fba SET status = 909, timestamp_updated = '{now}' WHERE fba.id NOT IN (SELECT DISTINCT id FROM {schema}.fb_sdk_adaccount) AND fba.status IS NOT NULL AND fba.status != 907 AND fba.fb_user_id IS NOT NULL; """ run_query(query_sql=query) return "updated" @staticmethod def list_adaccounts(schema) -> pd.DataFrame: """Getting the basic adaccounts information""" """907 - status is placed when adaccounts are being deleted""" query = f""" SELECT id, name, description, status, business_id, business_name FROM {schema}.fb_adaccount WHERE status != 907; """ df = run_query(query_sql=query, return_type="df") if df is None: df = pd.DataFrame( { "id": [], "name": [], "description": [], "status": [], "business_id": [], "business_name": [], } ) return df @staticmethod def mark_adaccounts_for_deletion(schema, fb_user_id): """Changing the status for adaccounts that are going to be deleted""" params = {"fb_user_id": str(fb_user_id)} query = f""" UPDATE {schema}.fb_adaccount SET status = 907 WHERE fb_user_id = %(fb_user_id)s; """ run_query(query, params) @staticmethod def delete_adaccount( schema, adaccount_id ): # !TODO Need to handle removal of shared audiences as well params = {"adaccount_id": adaccount_id} query = f""" DELETE FROM {schema}.fb_adaccount WHERE id = %(adaccount_id)s; """ run_query(query, params) @staticmethod def get_adaccounts_by_user_id(schema, user_id): """Should only be used with user_id internally obtained""" query = f""" SELECT * FROM {schema}.fb_adaccount WHERE fb_user_id = '{user_id}';""" df = run_query(query, return_type="df") return df @staticmethod def populate_fb_audience_shared_state_table(schema, internal_audience_id, data): for adaccount in data: now = pd.Timestamp.now() params = { "audience_id": internal_audience_id, "adaccount_id": adaccount["ad_acct_id"], "shared_status": adaccount["audience_share_status"], "timestamp_created": now, "timestamp_last_updated": now, } insert_fields = '("' + '", "'.join([key for key in params]) + '")' insert_values = "(" + ", ".join(["%(" + key + ")s" for key in params]) + ")" query = f""" INSERT INTO {schema}.fb_audience_shared_state {insert_fields} VALUES {insert_values};""" run_query(query, params) return data[0]["audience_share_status"] @staticmethod def delete_from_fb_audience_shared_state_table( adaccount_id, schema=None, management_schema=None, internal_audience_id=None, everything=False, ): """Removing shared state from audiences-adaccounts""" params = {"adaccount_id": adaccount_id, "audience_id": internal_audience_id} if schema is not None: schemas = [schema] elif management_schema is not None: schemas = list_caw_schemas(management_schema) else: raise RuntimeError("Invalid combination of parameters") if internal_audience_id is None and everything: where_query = """WHERE adaccount_id = %(adaccount_id)s""" else: where_query = """WHERE audience_id = %(audience_id)s AND adaccount_id = %(adaccount_id)s""" df_list = [] for schema in schemas: query = f""" DELETE FROM {schema}.fb_audience_shared_state {where_query} RETURNING *; """ df = run_query(query, params, return_type="df") df_list.append(df) result_df = pd.concat(df_list, ignore_index=True) return result_df @staticmethod def list_shared_audiences( adaccount_id, management_schema=None, external_id=False ) -> pd.DataFrame: """alliance_shcemas are provided when we work with multiple alliances""" params = {"adaccount_id": adaccount_id} external_id_part = "" if external_id: external_id_part = ",fa.external_id external_id" df_list = [] for schema in list_caw_schemas(management_schema): query = f""" SELECT fass.audience_id id, fa.name, fa.description, fa.actual_count size, collection_id {external_id_part} FROM {schema}.fb_audience_shared_state fass LEFT JOIN {schema}.fb_audience fa on fass.audience_id = fa.id WHERE adaccount_id = %(adaccount_id)s """ df = run_query(query, params, return_type="df") if df is None: df = pd.DataFrame( { "id": [], "name": [], "description": [], "size": [], "collection_id": [], } ) df_list.append(df) result_df = pd.concat(df_list, ignore_index=True) # Returning [Audience] appsync type return result_df @staticmethod def list_users(schema) -> pd.DataFrame: query = f""" SELECT user_id FROM {schema}.fb_user WHERE status != 'deleting'; """ df = run_query(query, return_type="df") return df @staticmethod def get_lookalike_ratios() -> pd.DataFrame: query = """SELECT id, name from commons.fb_lookalike_ratio;""" df = run_query(query, return_type="df") return df @staticmethod def get_countries() -> pd.DataFrame: query = """SELECT id, name from commons.fb_countries ORDER BY name;""" df = run_query(query, return_type="df") return df def populate_fb_user_table(self, schema, response, status): """Converting the obtained data and storing it in fb_user table""" df = self.convert_fb_account_response_to_df(response) df = df.iloc[0] # Ensuring that only one row is provided expires_in = df["expires_in"] df["status"] = status now = pd.Timestamp.now() df["token_expiration_timestamp"] = now + pd.to_timedelta(expires_in, unit="s") df["token_type"] = df.to_dict().get("token_type", "initial") df["timestamp_created"] = now df["timestamp_last_updated"] = now params = df.to_dict() insert_fields = '("' + '", "'.join([key for key in params]) + '")' insert_values = "(" + ", ".join(["%(" + key + ")s" for key in params]) + ")" on_conflict_values = ", ".join( [f"{key} = excluded.{key}" for key in params if key != "user_id"] ) query = f""" INSERT INTO {schema}.fb_user {insert_fields} VALUES {insert_values} ON CONFLICT ON CONSTRAINT unique_fb_user_id DO UPDATE SET {on_conflict_values} RETURNING user_id, status, granted_permissions; """ df = run_query(query, params, return_type="df") return df @staticmethod def convert_fb_account_response_to_df(obj): """Converts the received object from the front end into correct structure to insert in to db""" df = pd.DataFrame(obj, index=[0]) changed_columns = { "accessToken": "access_token", "data_access_expiration_time": "data_access_expiration_int", "expiresIn": "expires_in", "grantedScopes": "granted_permissions", "graphDomain": "graph_domain", "signedRequest": "signed_request", "userID": "user_id", } df.rename(columns=changed_columns, inplace=True) return df @staticmethod def update_fb_user_permissions(schema, user_id, permissions: List[str]): """Updating the granted_scopes field for the given user""" params = { "granted_permissions": ",".join(permissions), "user_id": user_id, "now": pd.Timestamp.now(), } query = f""" UPDATE {schema}.fb_user SET granted_permissions = %(granted_permissions)s, timestamp_last_updated = %(now)s WHERE user_id = %(user_id)s RETURNING user_id, granted_permissions; """ df = run_query(query, params, return_type="df") return df @staticmethod def delete_fb_user(schema): query = f""" DELETE FROM {schema}.fb_user; """ run_query(query) return "disconnected" @staticmethod def mark_fb_user_for_deletion(schema): """Marking fb_user that will be deleted""" query = f""" UPDATE {schema}.fb_user SET status = 'deleting'; """ run_query(query) @staticmethod def get_fb_user(schema): """Returning list of fb business manager accounts linked with FanSifter""" query = f""" SELECT * FROM {schema}.fb_user WHERE status != 'deleting'; """ df = run_query(query, return_type="df") return df