import hashlib import re from collections import defaultdict from datetime import datetime, timezone from typing import Any, BinaryIO import pandas as pd import phonenumbers from snowflake.snowpark import Session from backend import queries from backend.dtos import ( CampaignStats, FailureBreakdown, FileUploadFilterOption, FileUploadResult, UploadStatusResult, ) INPUT_FILE_COLUMNS = [ "EMAIL_REQUIRED", "OPTIN_REQUIRED", "MAILING_LIST", "FILE_SOURCE_DESCRIPTION", "TERRITORY_2_DIGIT_ISO_REQUIRED", "LABEL_REQUIRED", "FIRST_NAME", "LAST_NAME", "MOBILE_PHONE", "BIRTHDATE_MM_DD_YYYY", "BIRTHDAY_MM_DD_US_ONLY", "GENDER", "ADDRESS_1", "CITY", "STATE", "COUNTRY_REGION", "POSTAL_CODE", "PREFERRED_LANGUAGE", "FACEBOOK_PAGE", "TWITTER_HANDLE", ] DSPS = [ "Amazon", "Apple", "Deezer", "Email", "Facebook", "Google", "SMS", "Spotify", "Twitter", "YouTube", ] ACQUISITION_CHANNELS = { "B2B List": "EMAIL_REQUIRED", "Bandsintown List": "EMAIL_REQUIRED", "Event Shop": "EMAIL_REQUIRED", "Feature.FM": "EMAIL_REQUIRED", "Internal Sony List Transfer": "EMAIL_REQUIRED", "Levellr": "EMAIL_REQUIRED", "Linkfire": "EMAIL_REQUIRED", "Maestro List": "EMAIL_REQUIRED", "Medallion": "EMAIL_REQUIRED", "MTL": "EMAIL_REQUIRED", "Other": "EMAIL_REQUIRED", "Past Purchasers": "EMAIL_REQUIRED", "Previous Mailing List": "EMAIL_REQUIRED", "Seated List": "EMAIL_REQUIRED", "Seed List": "EMAIL_REQUIRED", "SMF": "EMAIL_REQUIRED", "SMS Campaigns": "MOBILE_PHONE", "Spotify - Sony Music Now": "EMAIL_REQUIRED", "Sqwad": "EMAIL_REQUIRED", "Test Upload": "EMAIL_REQUIRED", "Ticketing List": "EMAIL_REQUIRED", "Tunespeak": "EMAIL_REQUIRED", "Website": "EMAIL_REQUIRED", "Wyng": "EMAIL_REQUIRED", } IS_CRM_GENERATED = ["No", "Yes"] _ALLOWED_GENDERS = {"female", "male", "prefer not to answer", "non-binary/other"} def filter_options_df(options: list[FileUploadFilterOption]) -> pd.DataFrame: return pd.DataFrame([option.model_dump() for option in options]) def normalize_columns(dataframe: pd.DataFrame) -> pd.DataFrame: dataframe = dataframe.copy() dataframe.columns = dataframe.columns.str.strip() dataframe.columns = ( dataframe.columns.str.replace("(", "_", regex=False) .str.replace(")", "", regex=False) .str.replace("-", "", regex=False) .str.replace("/", "_", regex=False) .str.replace(" ", "_", regex=False) .str.replace("__", "_", regex=False) ) dataframe.columns = dataframe.columns.str.upper() return dataframe def validate_upload_columns(dataframe: pd.DataFrame) -> list[str]: if set(dataframe.columns) != set(INPUT_FILE_COLUMNS): return [ "Uploaded CSV file contains wrong column names or some columns are " "missing or extra." ] if dataframe["OPTIN_REQUIRED"].isna().sum() > 0: return ["Uploaded CSV file cannot contain missing values for opt-in column."] return [] def build_campaign_id( label_name: str, artist_name: str, file_source_description: str ) -> str: combined = f"{label_name}{artist_name}{file_source_description}" return hashlib.sha256(combined.encode()).hexdigest() def check_data_quality(dataframe: pd.DataFrame, upload_focus: str) -> list[str]: warnings: list[str] = [] total_rows = len(dataframe) missing_count = dataframe[upload_focus].isna().sum() duplicate_count = dataframe.duplicated(subset=upload_focus, keep="first").sum() if duplicate_count > 0: warnings.append( f"Warning! Number of duplicate {upload_focus} values: {duplicate_count}." ) if missing_count > 0: warnings.append( f"Warning! Column {upload_focus} is missing for {missing_count} rows " f"out of {total_rows} rows." ) if missing_count == total_rows: warnings.append( f"Error! Acquisition channel requires at least one row filled for " f"column {upload_focus}." ) return warnings # Ports the row-level rules from FANSIFTER_APP_REPORTING's # VALIDATE_FANRESPONSE_FROM_FILE_UPLOAD Snowflake function so bad rows can be # excluded and reported before they ever reach the DB-side validation task. _EMAIL_REGEX = re.compile( r"^[a-z0-9!#$%&'*+/=?^_`{|}~-]+(\.?|[a-z0-9!#$%&'*+/=?^_`{|}~-]+)*" r"@([a-z0-9]([a-z0-9-]*[a-z0-9])?\.)+[a-z0-9]([a-z0-9-]*[a-z0-9])?$" ) _ZERO_WIDTH_SPACE = "​" _DAYS_IN_MONTH = { 1: 31, 2: 29, 3: 31, 4: 30, 5: 31, 6: 30, 7: 31, 8: 31, 9: 30, 10: 31, 11: 30, 12: 31, } _LENGTH_LIMITS = [ # Snowflake's rule checks RECORD_CONTENT:address, but the record is built # from address_1 — validating ADDRESS_1 here matches the evident intent. ("ADDRESS_1", 255, "address"), ("CITY", 75, "city"), ("STATE", 255, "state"), ("COUNTRY_REGION", 56, "country_region"), ("POSTAL_CODE", 20, "postal_code"), ("PREFERRED_LANGUAGE", 255, "preferred_language"), ("FACEBOOK_PAGE", 255, "facebook_page"), ("TWITTER_HANDLE", 100, "twitter_handle"), ] _REQUIRED_FIELD_MESSAGES = [ ("FILE_SOURCE_DESCRIPTION", "Empty file_source_description"), ("ARTIST", "Empty artist"), ("TERRITORY", "Empty territory"), ("CAMPAIGN_ID", "Empty campaign_id"), ("CAMPAIGN_NAME", "Empty campaign_name"), ] def _is_present(value: Any) -> bool: return pd.notna(value) and str(value).strip() != "" def _is_valid_email_address(value: str) -> bool: cleaned = value.replace(_ZERO_WIDTH_SPACE, "").lower() return bool(_EMAIL_REGEX.match(cleaned)) def _is_valid_mobile_phone(value: str, country_region: object) -> bool: region = str(country_region) if _is_present(country_region) else None try: parsed = phonenumbers.parse(value, region=region) except phonenumbers.NumberParseException: return False return phonenumbers.is_valid_number(parsed) def _is_valid_birthdate_mm_dd_yyyy(value: str) -> bool: try: datetime.strptime(value, "%m/%d/%Y") except ValueError: return False return True def _is_valid_birthday_mm_dd(value: str) -> bool: parts = value.split("/") try: month = int(parts[0]) day = int(parts[1]) except (IndexError, ValueError): return False if not (1 <= month <= 12 and 1 <= day <= 31): return False return day <= _DAYS_IN_MONTH[month] def _is_valid_date_created(value: str) -> bool: try: datetime.fromisoformat(value) except ValueError: return False return True def _fanresponse_row_errors(row: pd.Series) -> list[str]: errors: list[str] = [] email = row.get("EMAIL_REQUIRED") mobile = row.get("MOBILE_PHONE") has_email = _is_present(email) has_mobile = _is_present(mobile) if has_email and not _is_valid_email_address(str(email)): errors.append("Invalid Email") if not has_email and not has_mobile: errors.append("Either email OR mobile_phone should be present") first_name = row.get("FIRST_NAME") if _is_present(first_name) and len(str(first_name)) > 110: errors.append("first_name is too long") last_name = row.get("LAST_NAME") if _is_present(last_name) and len(str(last_name)) > 110: errors.append("last_name is too long") if has_mobile and not _is_valid_mobile_phone( str(mobile), row.get("COUNTRY_REGION") ): errors.append("mobile_phone is invalid") birthdate = row.get("BIRTHDATE_MM_DD_YYYY") if _is_present(birthdate) and not _is_valid_birthdate_mm_dd_yyyy(str(birthdate)): errors.append("birthdate_mm_dd_yyyy is invalid") birthday = row.get("BIRTHDAY_MM_DD_US_ONLY") if _is_present(birthday) and not _is_valid_birthday_mm_dd(str(birthday)): errors.append("birthday_mm_dd_us_only is invalid") gender = row.get("GENDER") if _is_present(gender) and str(gender).lower() not in _ALLOWED_GENDERS: errors.append("Invalid gender") for column, max_length, label in _LENGTH_LIMITS: value = row.get(column) if _is_present(value) and len(str(value)) > max_length: errors.append(f"{label} is too long") date_created = row.get("DATE_CREATED") if not _is_present(date_created): errors.append("Empty date_created") elif not _is_valid_date_created(str(date_created)): errors.append("date_created is invalid") crm_generated = row.get("CRM_GENERATED") if not _is_present(crm_generated) or str(crm_generated).lower() not in ( "no", "yes", ): errors.append("crm_generated is invalid") dsp = row.get("DSP") if _is_present(dsp) and len(str(dsp)) > 255: errors.append("dsp is invalid") for column, message in _REQUIRED_FIELD_MESSAGES: if not _is_present(row.get(column)): errors.append(message) if row.get("OVERWRITE_FAN_DATA") not in (True, False): errors.append("overwrite_fan_data is invalid") return errors def validate_fan_data( dataframe: pd.DataFrame, ) -> tuple[pd.DataFrame, list[FailureBreakdown]]: if dataframe.empty: return dataframe, [] row_errors = dataframe.apply(_fanresponse_row_errors, axis=1) valid_mask = row_errors.apply(len) == 0 reason_counts: dict[str, int] = defaultdict(int) for errors in row_errors: for reason in errors: reason_counts[reason] += 1 failures = [ FailureBreakdown(reason=reason, count=count) for reason, count in sorted(reason_counts.items()) ] valid_dataframe = dataframe[valid_mask].reset_index(drop=True) return valid_dataframe, failures def prepare_upload_dataframe( file: BinaryIO, file_name: str, label_name: str, artist_name: str, mailing_list_name: str, file_source_description: str, dsp: str, is_crm_generated: str, acquisition_channel: str, overwrite_fan_data: bool, user: str, vendor_id: str | None = None, global_participant_id: str | None = None, virtual_participant_id: str | None = None, ) -> tuple[pd.DataFrame, list[str]]: dataframe = pd.read_csv(file, delimiter=",", dtype=defaultdict(lambda: str)) dataframe = normalize_columns(dataframe) errors = validate_upload_columns(dataframe) if errors: return dataframe, errors dataframe = dataframe.drop( ["LABEL_REQUIRED", "TERRITORY_2_DIGIT_ISO_REQUIRED"], axis=1 ) now = datetime.now(timezone.utc) dataframe["DATE_CREATED"] = now.isoformat() dataframe["CAMPAIGN_CREATED_BY"] = user dataframe["FILE_NAME"] = file_name dataframe["FILE_SOURCE_DESCRIPTION"] = file_source_description dataframe["LABEL"] = label_name dataframe["ARTIST"] = artist_name dataframe["MAILING_LIST"] = mailing_list_name dataframe["TLA_ID"] = None dataframe["VENDOR_ID"] = vendor_id dataframe["GLOBAL_PARTICIPANT_ID"] = global_participant_id dataframe["VIRTUAL_PARTICIPANT_ID"] = virtual_participant_id dataframe["ACQUISITION_CHANNEL"] = acquisition_channel dataframe["CRM_GENERATED"] = is_crm_generated dataframe["DSP"] = dsp dataframe["OVERWRITE_FAN_DATA"] = overwrite_fan_data dataframe["CAMPAIGN_NAME"] = ( f"{artist_name} - {file_source_description} - {now.strftime('%m/%d/%Y')}" ) dataframe["CAMPAIGN_ID"] = build_campaign_id( label_name, artist_name, file_source_description ) dataframe["GENDER"] = dataframe["GENDER"].astype(str) dataframe["GENDER"] = dataframe["GENDER"].apply( lambda g: g if g.lower() in _ALLOWED_GENDERS else pd.NA ) return dataframe, [] def submit_upload( session: Session, env: str, dataframe: pd.DataFrame, file_name: str ) -> FileUploadResult: if dataframe.empty: return FileUploadResult( success=False, campaign_id=None, rows_uploaded=0, errors=["No valid rows to upload."], ) existing_filenames = queries.query_existing_file_upload_filenames(session, env) if file_name in existing_filenames: return FileUploadResult( success=False, campaign_id=None, rows_uploaded=0, errors=["A file with this name has already been uploaded."], ) campaign_id = str(dataframe["CAMPAIGN_ID"].iloc[0]) existing_campaign_ids = queries.query_existing_campaign_ids(session, env) if campaign_id in existing_campaign_ids: return FileUploadResult( success=False, campaign_id=campaign_id, rows_uploaded=0, errors=[ "This campaign already exists, please update the input values to " "create a new campaign." ], ) rows_uploaded = queries.insert_file_upload_rows(session, dataframe, env) return FileUploadResult( success=True, campaign_id=campaign_id, rows_uploaded=rows_uploaded, errors=[] ) def check_upload_status( session: Session, env: str, upload_name: str ) -> UploadStatusResult: return queries.query_file_upload_status(session, upload_name, env) def list_campaign_artists(session: Session, env: str) -> list[str]: return queries.query_file_upload_campaign_artists(session, env) def list_campaign_stats( session: Session, env: str, artist_name: str ) -> list[CampaignStats]: return queries.query_file_upload_campaign_stats(session, artist_name, env) def delete_file_upload( session: Session, env: str, file_name: str, reason: str, username: str ) -> str: return queries.call_delete_file_upload_procedure( session, file_name, reason, username, env )