from typing import Optional import pandas as pd from snowflake.connector.errors import DatabaseError, ProgrammingError from snowflake.snowpark import Session from snowflake.snowpark.functions import call_udf, col, count from backend.dtos import ( CampaignStats, CcpaRow, DsarResult, ExportRecord, FailureBreakdown, FileUploadFilterOption, Profile, Subscription, UploadStatusResult, ) def query_fan_profile( session: Session, email: str, env: str = "QA" ) -> Optional[Profile]: escaped = email.replace("'", "''") deleted_rows = session.sql( f"SELECT EMAIL, DELETED_AT FROM PREFERENCE_CENTER.{env}.DELETED_PROFILE " f"WHERE lower(EMAIL) = lower('{escaped}')" ).collect() deleted_at = deleted_rows[0]["DELETED_AT"] if deleted_rows else None results = session.sql( f"SELECT ID, CRM_ID, EMAIL, FIRST_NAME, LAST_NAME, PHONE_NUMBER, COUNTRY_CODE, " f"CITY, REGION, ADDRESS, ZIP_CODE, DATE_OF_BIRTH, BIRTHDAY, DELETED " f"FROM PREFERENCE_CENTER.{env}.FAN_PROFILE " f"WHERE lower(EMAIL) = lower('{escaped}')" ).collect() if not results: if deleted_at is not None: return Profile( profile_id=None, crm_id=None, email=email, first_name=None, last_name=None, phone_number=None, country_code=None, city=None, region=None, address=None, zip_code=None, date_of_birth=None, birthday=None, deleted=True, deleted_at=deleted_at, ) return None row = results[0] return Profile( profile_id=row["ID"], crm_id=row["CRM_ID"], email=row["EMAIL"], first_name=row["FIRST_NAME"], last_name=row["LAST_NAME"], phone_number=row["PHONE_NUMBER"], country_code=row["COUNTRY_CODE"], city=row["CITY"], region=row["REGION"], address=row["ADDRESS"], zip_code=row["ZIP_CODE"], date_of_birth=row["DATE_OF_BIRTH"], birthday=row["BIRTHDAY"], deleted=bool(row["DELETED"]), deleted_at=deleted_at, ) def query_fan_subscriptions( session: Session, email: str, env: str = "QA" ) -> list[Subscription]: escaped = email.replace("'", "''") rows = session.sql( f"SELECT fs.CREATED_AT, ml.NAME AS ARTIST_NAME " f"FROM PREFERENCE_CENTER.{env}.FAN_SUBSCRIPTION fs " f"INNER JOIN PREFERENCE_CENTER.{env}.FAN_PROFILE fp ON fs.PROFILE_ID = fp.ID " f"INNER JOIN PREFERENCE_CENTER.{env}.FAN_MAILING_LIST ml ON fs.MAILING_LIST_ID = ml.ID " f"WHERE lower(fp.EMAIL) = lower('{escaped}') AND fs.IS_ACTIVE = TRUE " f"ORDER BY fs.CREATED_AT" ).collect() return [ Subscription(created_at=r["CREATED_AT"], artist_name=r["ARTIST_NAME"]) for r in rows ] def query_fan_exports( session: Session, email: str, env: str = "QA" ) -> tuple[list[ExportRecord], list[str]]: escaped = email.lower().replace("'", "''") exports: list[ExportRecord] = [] warnings: list[str] = [] try: fansifter_rows = session.sql( f"SELECT distinct ex.reason, ex.snapshot_id, id.name as exporter_name, " f"pro.first_name, pro.email, ex.justification, " f"max(ex.created_at) as last_export_date " f"FROM fansifter_pg_reporting.{env.lower()}_ows_dmp_public.audience_export ex " f"INNER JOIN fansifter_app_reporting.{env}.audience_fan snap ON ex.snapshot_id=snap.snapshot_id " f"INNER JOIN preference_center.{env}.fan_profile pro ON snap.fan_id=sha2(lower(pro.email), 256) " f"LEFT JOIN FACTS.{env}.IDENTITY id ON ex.created_by=id.id " f"WHERE ex.status='COMPLETED' AND NOT ex._fivetran_deleted " f"AND NOT pro.deleted " f"AND snap.fan_id=sha2(lower('{escaped}'),256) " f"GROUP BY all" ).collect() except Exception as e: fansifter_rows = [] warnings.append(f"Fansifter export data could not be loaded: {e}") try: salesforce_rows = session.sql( f"SELECT eh.created_date AS export_date, " f"e.business_justification_c AS business_reason, " f"e.recipient_email_c, e.recipient_name_c, e.export_type_c, e.exported_by_c, " f"eh.first_name_c, eh.last_name_c, eh.email_c " f"FROM delphi_crm_data.raw_salesforce_sales_cloud.export_c e " f"INNER JOIN delphi_crm_data.raw_salesforce_sales_cloud.export_history_c eh " f"ON eh.export_id_c = e.id " f"WHERE lower(eh.email_c) = lower('{escaped}')" ).collect() except Exception as e: salesforce_rows = [] warnings.append(f"Salesforce export data could not be loaded: {e}") for r in fansifter_rows: exports.append( ExportRecord( export_date=r["LAST_EXPORT_DATE"], business_reason=r["JUSTIFICATION"], recipient_email=None, recipient_name=r["EXPORTER_NAME"], export_type=r["REASON"], exported_by=r["EXPORTER_NAME"], fan_first_name=r["FIRST_NAME"], fan_last_name=None, fan_email=r["EMAIL"], source="Fansifter", ) ) try: hightouch_rows = session.sql( f"SELECT sr.model_name, sr.destination, sr.started_at AS export_date, cl.op_type, " f"cl.fields:EMAIL_C::string AS fan_email " f"FROM delphi_crm_data.hightouch_audit.sync_changelog cl " f"INNER JOIN delphi_crm_data.hightouch_audit.sync_runs sr " f"ON sr.sync_id = cl.sync_id AND sr.sync_run_id = cl.sync_run_id " f"WHERE lower(cl.fields:EMAIL_C::string) = '{escaped}' " f"AND cl.status = 'succeeded' " f"ORDER BY sr.started_at DESC" ).collect() except Exception as e: hightouch_rows = [] warnings.append(f"Hightouch export data could not be loaded: {e}") for r in salesforce_rows: exports.append( ExportRecord( export_date=r["EXPORT_DATE"], business_reason=r["BUSINESS_REASON"], recipient_email=r["RECIPIENT_EMAIL_C"], recipient_name=r["RECIPIENT_NAME_C"], export_type=r["EXPORT_TYPE_C"], exported_by=r["EXPORTED_BY_C"], fan_first_name=r["FIRST_NAME_C"], fan_last_name=r["LAST_NAME_C"], fan_email=r["EMAIL_C"], source="Salesforce", ) ) for r in hightouch_rows: destination = r["DESTINATION"] or "" business_reason = ( "email audience" if destination.lower() == "sfmc" else "retargeting" ) exports.append( ExportRecord( export_date=r["EXPORT_DATE"], business_reason=business_reason, recipient_email=None, recipient_name=destination, export_type=r["MODEL_NAME"], exported_by=None, fan_first_name=None, fan_last_name=None, fan_email=r["FAN_EMAIL"], source="Hightouch", ) ) return exports, warnings def query_dsar_report(session: Session, email: str, env: str = "QA") -> DsarResult: profile = query_fan_profile(session, email, env) subscriptions = query_fan_subscriptions(session, email, env) exports, export_warnings = query_fan_exports(session, email, env) return DsarResult( profile=profile, subscriptions=subscriptions, exports=exports, export_warnings=export_warnings, ) def upsert_ccpa_emails( session: Session, rows: "list[CcpaRow]", user: str, env: str = "QA", ) -> tuple[int, int, int]: if not rows: return (0, 0, 0) table_name = f"FANSIFTER_APP_REPORTING.{env}.CCPA_FAN_LIST" escaped_user = user.replace("'", "''") values_list = ", ".join( [ f"('{row.email.replace(chr(39), chr(39) + chr(39))}', '{row.requested_date.isoformat()}')" for row in rows ] ) warehouse = session.get_current_warehouse() if warehouse: session.sql(f"USE WAREHOUSE {warehouse.strip(chr(34))}").collect() merge_sql = f""" MERGE INTO {table_name} AS target USING ( SELECT column1 AS EMAIL, column2::DATE AS REQUESTED_AT FROM VALUES {values_list} ) AS source ON target.EMAIL = source.EMAIL WHEN MATCHED AND source.REQUESTED_AT < target.REQUESTED_AT THEN UPDATE SET REQUESTED_AT = source.REQUESTED_AT, UPDATED_AT = CONVERT_TIMEZONE('UTC', CURRENT_TIMESTAMP())::TIMESTAMP_NTZ, UPDATED_BY = '{escaped_user}' WHEN NOT MATCHED THEN INSERT (EMAIL, REQUESTED_AT, CREATED_AT, UPDATED_AT, CREATED_BY, UPDATED_BY) VALUES (source.EMAIL, source.REQUESTED_AT, CONVERT_TIMEZONE('UTC', CURRENT_TIMESTAMP())::TIMESTAMP_NTZ, CONVERT_TIMEZONE('UTC', CURRENT_TIMESTAMP())::TIMESTAMP_NTZ, '{escaped_user}', '{escaped_user}') """ result = session.sql(merge_sql).collect() if result: row = result[0] inserted = int(row["number of rows inserted"]) updated = int(row["number of rows updated"]) skipped = len(rows) - inserted - updated return (inserted, updated, skipped) return (0, 0, len(rows)) def insert_deleted_profiles( session: Session, emails: list[str], user: str, env: str = "QA", ) -> int: if not emails: return 0 from snowflake.snowpark import functions as F fan_profile_table = f"PREFERENCE_CENTER.{env}.FAN_PROFILE" deleted_profile_table = f"PREFERENCE_CENTER.{env}.DELETED_PROFILE" warehouse = session.get_current_warehouse() if warehouse: session.sql(f"USE WAREHOUSE {warehouse.strip(chr(34))}").collect() emails_df = session.create_dataframe([[e] for e in emails], schema=["INPUT_EMAIL"]) fan_profile_df = session.table(fan_profile_table).select( F.col("ID").alias("PROFILE_ID"), F.col("EMAIL").alias("PROFILE_EMAIL"), ) matched_df = fan_profile_df.join( emails_df, F.lower("PROFILE_EMAIL") == F.lower("INPUT_EMAIL"), join_type="inner", ).select( F.col("PROFILE_ID"), F.col("PROFILE_EMAIL").alias("EMAIL"), F.sql_expr("CONVERT_TIMEZONE('UTC', CURRENT_TIMESTAMP())::TIMESTAMP_NTZ").alias( "DELETED_AT" ), F.lit("BULK_DELETE_EVENT").alias("REASON"), F.lit(user).alias("UPDATED_BY"), ) matched_df.write.save_as_table( deleted_profile_table, mode="append", column_order="name" ) return int(matched_df.count()) def query_file_upload_filter_options( session: Session, env: str = "QA" ) -> list[FileUploadFilterOption]: rows = session.sql( f""" with fansifter_artist as ( select vendor_id, global_participant_id as id, 'artist' as type from fansifter_app_reporting.{env}.artist_roster_main_rep union select vendor_id, global_participant_id as id, 'artist' as type from fansifter_app_reporting.{env}.artist_roster_local_rep union select vendor_id, global_participant_id as id, 'artist' as type from fansifter_app_reporting.{env}.artist_roster ) select fa.vendor_id, v.name as label_name, fa.id as global_participant_id, NULL as virtual_participant_id, gp.name as artist_name, ml.id as mailing_list_id, ml.name as mailing_list_name, fa.type as type from fansifter_artist fa inner join facts.{env}.vendor v on fa.vendor_id = v.vendor_id inner join facts.{env}.global_participant gp on gp.id = fa.id inner join PREFERENCE_CENTER.{env}.FAN_MAILING_LIST ML ON gp.id=ml.global_participant_id where not v.is_deleted union all select cl.vendor_id, v.name as label_name, NULL as global_participant_id, vp.id as virtual_participant_id, vp.name as artist_name, ml.id as mailing_list_id, ml.name as mailing_list_name, 'custom list' as type from PREFERENCE_CENTER.{env}.FAN_MAILING_LIST ML inner join fansifter_app_reporting.{env}.virtual_participant vp on ml.virtual_participant_id=vp.id inner join fansifter_app_reporting.{env}.custom_list cl on cl.virtual_participant_id=vp.id inner join facts.{env}.vendor v on cl.vendor_id = v.vendor_id where not v.is_deleted """ ).collect() return [ FileUploadFilterOption( vendor_id=str(r["VENDOR_ID"]), label_name=r["LABEL_NAME"], global_participant_id=( str(r["GLOBAL_PARTICIPANT_ID"]) if r["GLOBAL_PARTICIPANT_ID"] is not None else None ), virtual_participant_id=( str(r["VIRTUAL_PARTICIPANT_ID"]) if r["VIRTUAL_PARTICIPANT_ID"] is not None else None ), artist_name=r["ARTIST_NAME"], mailing_list_id=str(r["MAILING_LIST_ID"]), mailing_list_name=r["MAILING_LIST_NAME"], type=r["TYPE"], ) for r in rows ] def _crm_fans_schema(env: str) -> str: return f"FANSIFTER_APP_REPORTING.{env}_CRM_FANS" def query_existing_file_upload_filenames( session: Session, env: str = "QA" ) -> list[str]: try: rows = ( session.table(f"{_crm_fans_schema(env)}.EVENT_FILE_UPLOAD") .select(col("FILE_NAME")) .distinct() .to_pandas() ) return [str(v) for v in rows.values.flatten().tolist()] except (ProgrammingError, DatabaseError, ValueError): return [] def query_existing_campaign_ids(session: Session, env: str = "QA") -> list[str]: try: rows = ( session.table(f"{_crm_fans_schema(env)}.EVENT_FILE_UPLOAD") .select(col("CAMPAIGN_ID")) .distinct() .to_pandas() ) return [str(v) for v in rows.values.flatten().tolist()] except (ProgrammingError, DatabaseError, ValueError): return [] def insert_file_upload_rows( session: Session, dataframe: pd.DataFrame, env: str = "QA" ) -> int: table_name = f"{_crm_fans_schema(env)}.EVENT_FILE_UPLOAD" session.create_dataframe(dataframe).write.save_as_table( table_name, mode="append", column_order="name" ) return len(dataframe) def call_delete_file_upload_procedure( session: Session, file_name: str, reason: str, username: str, env: str = "QA" ) -> str: return str( session.call( f"{_crm_fans_schema(env)}.DELETE_BACKOFFICE_FILE_UPLOAD", file_name, reason, username, ) ) def query_file_upload_status( session: Session, upload_name: str, env: str = "QA" ) -> UploadStatusResult: try: processed_count = ( session.table(f"{_crm_fans_schema(env)}.EVENT_FANRESPONSE_VALIDATED") .filter(col("FILE_SOURCE_DESCRIPTION") == upload_name) .select(col("EVENT_ROW_ID")) .distinct() .count() ) failure_rows = ( session.table(f"{_crm_fans_schema(env)}.EVENT_FANRESPONSE_ERROR") .filter(col("RECORD_CONTENT")["file_source_description"] == upload_name) .groupBy("NOT_VALID_REASON") .agg(count(col("EVENT_ROW_ID")).as_("FANS")) .collect() ) failures = [ FailureBreakdown(reason=r["NOT_VALID_REASON"], count=r["FANS"]) for r in failure_rows ] except (ProgrammingError, DatabaseError, ValueError): processed_count = 0 failures = [] return UploadStatusResult(processed_count=processed_count, failures=failures) def query_file_upload_campaign_artists(session: Session, env: str = "QA") -> list[str]: try: rows = ( session.table(f"{_crm_fans_schema(env)}.FILE_UPLOAD_CAMPAIGNS") .select(col("ARTIST_NAME")) .distinct() .to_pandas() ) return sorted(str(v) for v in rows.values.flatten().tolist()) except (ProgrammingError, DatabaseError, ValueError): return [] def query_file_upload_campaign_stats( session: Session, artist_name: str, env: str = "QA" ) -> list[CampaignStats]: escaped = artist_name.replace("'", "''") schema = _crm_fans_schema(env) try: rows = session.sql( f"SELECT uc.ARTIST_NAME, uc.FILE_SOURCE_DESCRIPTION, efu.FILE_NAME, " f"COUNT(*) AS FAN_COUNT " f"FROM {schema}.FILE_UPLOAD_CAMPAIGNS uc " f"INNER JOIN {schema}.EVENT_FANRESPONSE_VALIDATED efv " f"ON uc.CAMPAIGN_ID = efv.FORM_ID " f"INNER JOIN {schema}.EVENT_FILE_UPLOAD efu " f"ON efu.ROW_ID::VARCHAR = efv.EVENT_ROW_ID " f"WHERE uc.ARTIST_NAME = '{escaped}' " f"GROUP BY ALL" ).collect() except (ProgrammingError, DatabaseError, ValueError): return [] return [ CampaignStats( artist_name=r["ARTIST_NAME"], file_source_description=r["FILE_SOURCE_DESCRIPTION"], file_name=r["FILE_NAME"], fan_count=r["FAN_COUNT"], ) for r in rows ] def convert_country_to_iso2( session: Session, dataframe: pd.DataFrame, env: str = "QA" ) -> pd.DataFrame: snowpark_df = session.create_dataframe(dataframe) converted: pd.DataFrame = snowpark_df.withColumn( "COUNTRY_REGION", call_udf(f"{_crm_fans_schema(env)}.COUNTRY_TO_ISO2", col("COUNTRY_REGION")), ).to_pandas() return converted