"""Business logic services for artist roster operations. Environment Safety: - Local development: target=qa (enforced by Makefile) - PROD environment: Only accessible via Snowflake deployment - Write operations blocked when IS_PROD == True for extra safety """ from typing import Any import pandas as pd from snowflake.snowpark.session import Session from common import queries from common.db import execute_non_query, execute_query, get_current_user from common.environment import IS_PROD from common.types import ( AddArtistForm, ArtistRow, RosterFilters, UpgradePendingArtistForm, VendorRow, VendorWithBrandRow, VirtualParticipantForm, VirtualParticipantRow, ) from common.utils import parse_spotify_id def search_vendors( session: Session, vendor_id: int | None = None, vendor_name: str | None = None ) -> list[VendorRow]: """ Search for vendors by ID or name. Args: session: Active Snowpark session vendor_id: Optional vendor ID for exact match vendor_name: Optional vendor_name for pattern match Returns: List of VendorRow objects """ where_conditions = [] if vendor_id: where_conditions.append(f"VENDOR_ID = {vendor_id}") if vendor_name: escaped_name = vendor_name.replace("'", "''") where_conditions.append(f"NAME ILIKE '%{escaped_name}%'") where_clause = " OR ".join(where_conditions) if where_conditions else "1=1" query = f""" SELECT DISTINCT VENDOR_ID, NAME FROM orchard_app_reporting.delphi_prod.VENDOR WHERE NOT _FIVETRAN_DELETED AND ({where_clause}) ORDER BY NAME LIMIT 50 """ df = execute_query(session, query) if df.empty: return [] return [VendorRow(**row) for row in df.to_dict("records")] # type: ignore[misc] def search_artists( session: Session, artist_uuid: str | None = None, artist_name: str | None = None, spotify_id: str | None = None, ) -> list[ArtistRow]: """ Search for artists by UUID, name, or Spotify ID. Spotify ID supports multiple formats: - URL with params: https://open.spotify.com/artist/7tNO3vJC9zlHy2IJOx34ga?si=... - URL without params: https://open.spotify.com/artist/7tNO3vJC9zlHy2IJOx34ga - URI format: spotify:artist:4OGiMt96TFUKkKWf7Imlno - Plain ID: 7tNO3vJC9zlHy2IJOx34ga Args: session: Active Snowpark session artist_uuid: Optional artist UUID for exact match artist_name: Optional artist name for pattern match spotify_id: Optional Spotify ID (any format) for exact match Returns: List of ArtistRow objects """ where_conditions = [] if artist_uuid: where_conditions.append(f"ID = '{artist_uuid}'") if artist_name: escaped_name = artist_name.replace("'", "''") where_conditions.append(f"NAME ILIKE '%{escaped_name}%'") if spotify_id: # Parse Spotify ID from URL/URI/plain format parsed_id = parse_spotify_id(spotify_id) where_conditions.append(f"SPOTIFY_ID = '{parsed_id}'") where_clause = " OR ".join(where_conditions) if where_conditions else "1=1" query = f""" SELECT DISTINCT ID AS ARTIST_UUID, NAME AS ARTIST_NAME, SPOTIFY_ID FROM orchard_app_reporting.delphi_prod.GLOBAL_PARTICIPANT WHERE NOT _FIVETRAN_DELETED AND ({where_clause}) ORDER BY NAME LIMIT 50 """ df = execute_query(session, query) if df.empty: return [] return [ArtistRow(**row) for row in df.to_dict("records")] # type: ignore[misc] def get_all_roster_vendors(session: Session) -> list[VendorWithBrandRow]: """ Get all unique vendors from roster tables with their brands. Each vendor has exactly one brand. Returns list of vendors sorted by name. Args: session: Active Snowpark session Returns: List of VendorWithBrandRow objects sorted by vendor name """ schema_name = queries.SCHEMA_NAME query = f""" WITH all_roster_vendors AS ( SELECT DISTINCT VENDOR_ID FROM {schema_name}.ARTIST_ROSTER UNION SELECT DISTINCT VENDOR_ID FROM {schema_name}.ARTIST_ROSTER_MAIN_REP UNION SELECT DISTINCT VENDOR_ID FROM {schema_name}.ARTIST_ROSTER_LOCAL_REP ) SELECT v.VENDOR_ID as vendor_id, v.NAME as vendor_name, cb.DISPLAY_NAME as brand_name FROM all_roster_vendors arv JOIN orchard_app_reporting.delphi_prod.VENDOR v ON v.VENDOR_ID = arv.VENDOR_ID AND NOT v._FIVETRAN_DELETED LEFT JOIN orchard_app_reporting.delphi_prod.COMPANY_BRAND_HAS_LABEL_VENDOR cblv ON cblv.VENDOR_ID = v.VENDOR_ID AND NOT cblv._FIVETRAN_DELETED LEFT JOIN orchard_app_reporting.delphi_prod.COMPANY_BRAND cb ON cb.UUID = cblv.COMPANY_BRAND_UUID AND NOT cb._FIVETRAN_DELETED ORDER BY v.NAME """ df = execute_query(session, query) if df.empty: return [] return [VendorWithBrandRow(**row) for row in df.to_dict("records")] # type: ignore[misc] def is_vendor_non_sme(session: Session, vendor_id: int) -> bool: """ Check if a vendor is non-SME brand. DEPRECATED: Use is_vendor_sme() instead for accurate brand detection. A vendor is non-SME if it has at least one row in ARTIST_ROSTER table. Args: session: Active Snowpark session vendor_id: Vendor ID to check Returns: True if vendor is non-SME, False if SME """ schema_name = queries.SCHEMA_NAME query = f""" SELECT 1 AS is_non_sme FROM {schema_name}.ARTIST_ROSTER WHERE VENDOR_ID = {vendor_id} LIMIT 1 """ df = execute_query(session, query) return not df.empty def is_vendor_sme(session: Session, vendor_id: int) -> bool: """ Check if a vendor has SME (Sony Music) brand. Args: session: Active Snowpark session vendor_id: Vendor ID to check Returns: True if vendor has Sony Music brand, False otherwise """ query = f""" SELECT 1 AS is_sme FROM orchard_app_reporting.delphi_prod.VENDOR v JOIN orchard_app_reporting.delphi_prod.COMPANY_BRAND_HAS_LABEL_VENDOR cblv ON cblv.VENDOR_ID = v.VENDOR_ID AND NOT cblv._FIVETRAN_DELETED JOIN orchard_app_reporting.delphi_prod.COMPANY_BRAND cb ON cb.UUID = cblv.COMPANY_BRAND_UUID AND NOT cb._FIVETRAN_DELETED WHERE v.VENDOR_ID = {vendor_id} AND cb.DISPLAY_NAME = 'Sony Music' AND NOT v._FIVETRAN_DELETED LIMIT 1 """ df = execute_query(session, query) return not df.empty def check_sme_cross_roster_conflict( session: Session, vendor_id: int, artist_uuid: str, subaccount_id: int, target_roster: str, ) -> tuple[bool, str | None]: """ Check if artist exists in opposite SME roster. Business rule: SME artists cannot be in both MAIN_REP and LOCAL_REP. Args: session: Active Snowpark session vendor_id: Vendor ID artist_uuid: Artist UUID subaccount_id: Subaccount ID target_roster: Target roster type ('MAIN_REP' or 'LOCAL_REP') Returns: Tuple of (has_conflict: bool, existing_roster: str | None) """ schema_name = queries.SCHEMA_NAME query = f""" SELECT 'MAIN_REP' as existing_roster FROM {schema_name}.ARTIST_ROSTER_MAIN_REP WHERE VENDOR_ID = {vendor_id} AND GLOBAL_PARTICIPANT_ID = '{artist_uuid}' AND SUBACCOUNT_ID = {subaccount_id} UNION ALL SELECT 'LOCAL_REP' as existing_roster FROM {schema_name}.ARTIST_ROSTER_LOCAL_REP WHERE VENDOR_ID = {vendor_id} AND GLOBAL_PARTICIPANT_ID = '{artist_uuid}' AND SUBACCOUNT_ID = {subaccount_id} """ df = execute_query(session, query) if df.empty: return False, None # Get existing roster types (check for uppercase column name from Snowflake) col_name = ( "EXISTING_ROSTER" if "EXISTING_ROSTER" in df.columns else "existing_roster" ) existing_rosters = df[col_name].tolist() # Check for conflict with target roster if target_roster == "MAIN_REP" and "LOCAL_REP" in existing_rosters: return True, "LOCAL_REP" elif target_roster == "LOCAL_REP" and "MAIN_REP" in existing_rosters: return True, "MAIN_REP" return False, None def check_duplicate(session: Session, form: AddArtistForm) -> bool: """ Check if artist roster entry already exists. Args: session: Active Snowpark session form: Add artist form data Returns: True if duplicate exists, False otherwise """ if form.roster_type == "MAIN_REP": query = f""" SELECT 1 FROM {queries.SCHEMA_NAME}.ARTIST_ROSTER_MAIN_REP WHERE VENDOR_ID = {form.vendor_id} AND GLOBAL_PARTICIPANT_ID = '{form.artist_uuid}' AND SUBACCOUNT_ID = {form.subaccount_id} LIMIT 1 """ elif form.roster_type == "LOCAL_REP": query = f""" SELECT 1 FROM {queries.SCHEMA_NAME}.ARTIST_ROSTER_LOCAL_REP WHERE VENDOR_ID = {form.vendor_id} AND GLOBAL_PARTICIPANT_ID = '{form.artist_uuid}' AND SUBACCOUNT_ID = {form.subaccount_id} AND COUNTRY_CODE = '{form.country_code}' LIMIT 1 """ else: # ARTIST_ROSTER query = f""" SELECT 1 FROM {queries.SCHEMA_NAME}.ARTIST_ROSTER WHERE VENDOR_ID = {form.vendor_id} AND GLOBAL_PARTICIPANT_ID = '{form.artist_uuid}' AND SUBACCOUNT_ID = {form.subaccount_id} LIMIT 1 """ df = execute_query(session, query) return not df.empty def validate_add_artist(session: Session, form: AddArtistForm) -> list[str]: """ Validate add artist form data. Args: session: Active Snowpark session form: Add artist form data Returns: List of validation error messages (empty if valid) """ errors = form.validate_fields() # Check vendor exists vendors = search_vendors(session, vendor_id=form.vendor_id) if not vendors: errors.append(f"Vendor ID {form.vendor_id} does not exist") return errors # Stop validation if vendor doesn't exist # Check artist exists artists = search_artists(session, artist_uuid=form.artist_uuid) if not artists: errors.append(f"Artist UUID {form.artist_uuid} does not exist") return errors # Stop validation if artist doesn't exist # Brand detection and roster type validation is_sme = is_vendor_sme(session, form.vendor_id) if is_sme: # SME vendors can only use MAIN_REP or LOCAL_REP if form.roster_type not in ["MAIN_REP", "LOCAL_REP"]: errors.append( f"Vendor {form.vendor_id} is SME brand (Sony Music) and can only add to MAIN_REP or LOCAL_REP" ) # Check cross-roster conflict for SME vendors if form.roster_type in ["MAIN_REP", "LOCAL_REP"]: has_conflict, existing_roster = check_sme_cross_roster_conflict( session, form.vendor_id, form.artist_uuid, form.subaccount_id, form.roster_type, ) if has_conflict: errors.append( f"Artist already exists in {existing_roster} for this vendor. " f"Cannot add to {form.roster_type}. Use update functionality in the future." ) else: # Non-SME vendors can only use ARTIST_ROSTER if form.roster_type != "ARTIST_ROSTER": errors.append( f"Vendor {form.vendor_id} is non-SME brand and can only add to ARTIST_ROSTER" ) return errors def add_artist_to_roster(session: Session, form: AddArtistForm) -> None: """ Add artist to roster table using MERGE to prevent duplicates. Args: session: Active Snowpark session form: Add artist form data Raises: ValueError: If validation fails or PROD writes are blocked """ # CRITICAL: Block all PROD writes for safety in QA environment if IS_PROD: raise ValueError( "PROD writes are BLOCKED for safety. " "Use QA schema for testing. Set target=qa environment variable for QA writes." ) # Validate first errors = validate_add_artist(session, form) if errors: raise ValueError("; ".join(errors)) schema_name = queries.SCHEMA_NAME # Get current user for CREATED_BY tracking current_user = get_current_user(session) # Build appropriate MERGE query based on roster type if form.roster_type == "MAIN_REP": is_artist_team_val = "TRUE" if form.is_artist_team else "FALSE" query = f""" MERGE INTO {schema_name}.ARTIST_ROSTER_MAIN_REP t USING ( SELECT '{form.artist_uuid}' AS GLOBAL_PARTICIPANT_ID, {form.vendor_id} AS VENDOR_ID, {form.subaccount_id} AS SUBACCOUNT_ID, '{form.status}' AS STATUS, {is_artist_team_val} AS IS_ARTIST_TEAM, '{current_user}' AS CREATED_BY ) s ON ( t.GLOBAL_PARTICIPANT_ID = s.GLOBAL_PARTICIPANT_ID AND t.VENDOR_ID = s.VENDOR_ID AND t.SUBACCOUNT_ID = s.SUBACCOUNT_ID ) WHEN NOT MATCHED THEN INSERT (GLOBAL_PARTICIPANT_ID, VENDOR_ID, SUBACCOUNT_ID, STATUS, IS_ARTIST_TEAM, CREATED_BY) VALUES (s.GLOBAL_PARTICIPANT_ID, s.VENDOR_ID, s.SUBACCOUNT_ID, s.STATUS, s.IS_ARTIST_TEAM, s.CREATED_BY) """ elif form.roster_type == "LOCAL_REP": query = f""" MERGE INTO {schema_name}.ARTIST_ROSTER_LOCAL_REP t USING ( SELECT '{form.artist_uuid}' AS GLOBAL_PARTICIPANT_ID, {form.vendor_id} AS VENDOR_ID, {form.subaccount_id} AS SUBACCOUNT_ID, '{form.country_code}' AS COUNTRY_CODE, '{current_user}' AS CREATED_BY ) s ON ( t.GLOBAL_PARTICIPANT_ID = s.GLOBAL_PARTICIPANT_ID AND t.VENDOR_ID = s.VENDOR_ID AND t.SUBACCOUNT_ID = s.SUBACCOUNT_ID AND t.COUNTRY_CODE = s.COUNTRY_CODE ) WHEN NOT MATCHED THEN INSERT (GLOBAL_PARTICIPANT_ID, VENDOR_ID, SUBACCOUNT_ID, COUNTRY_CODE, CREATED_BY) VALUES (s.GLOBAL_PARTICIPANT_ID, s.VENDOR_ID, s.SUBACCOUNT_ID, s.COUNTRY_CODE, s.CREATED_BY) """ else: # ARTIST_ROSTER query = f""" MERGE INTO {schema_name}.ARTIST_ROSTER t USING ( SELECT '{form.artist_uuid}' AS GLOBAL_PARTICIPANT_ID, {form.vendor_id} AS VENDOR_ID, {form.subaccount_id} AS SUBACCOUNT_ID, '{current_user}' AS CREATED_BY ) s ON ( t.GLOBAL_PARTICIPANT_ID = s.GLOBAL_PARTICIPANT_ID AND t.VENDOR_ID = s.VENDOR_ID AND t.SUBACCOUNT_ID = s.SUBACCOUNT_ID ) WHEN NOT MATCHED THEN INSERT (GLOBAL_PARTICIPANT_ID, VENDOR_ID, SUBACCOUNT_ID, CREATED_BY) VALUES (s.GLOBAL_PARTICIPANT_ID, s.VENDOR_ID, s.SUBACCOUNT_ID, s.CREATED_BY) """ # Execute merge execute_non_query(session, query) def get_roster_view(session: Session, filters: RosterFilters) -> pd.DataFrame: """ Get roster view with filters. Args: session: Active Snowpark session filters: Search filters Returns: Pandas DataFrame with roster results """ # Build dynamic WHERE conditions where_conditions = [] if filters.vendor_id: where_conditions.append(f"a.VENDOR_ID = {filters.vendor_id}") if filters.vendor_name: escaped_name = filters.vendor_name.replace("'", "''") where_conditions.append(f"v.NAME ILIKE '%{escaped_name}%'") if filters.artist_uuid: where_conditions.append(f"a.artist_uuid = '{filters.artist_uuid}'") if filters.artist_name: escaped_name = filters.artist_name.replace("'", "''") where_conditions.append(f"gp.NAME ILIKE '%{escaped_name}%'") if filters.subaccount_id: where_conditions.append(f"a.SUBACCOUNT_ID = {filters.subaccount_id}") where_clause = " AND ".join(where_conditions) if where_conditions else "1=1" limit = filters.limit if filters.limit is not None else 100 offset = filters.offset if filters.offset is not None else 0 query = queries.ROSTER_VIEW_QUERY_TEMPLATE.format( where_clause=where_clause, limit=limit, offset=offset ) return execute_query(session, query) def delete_roster_entry( session: Session, vendor_id: int, artist_uuid: str, subaccount_id: int, roster_type: str, country_code: str | None = None, delete_all_countries: bool = False, ) -> None: """ Delete artist from roster table. Args: session: Active Snowpark session vendor_id: Vendor ID artist_uuid: Artist UUID subaccount_id: Subaccount ID roster_type: Roster type ('MAIN', 'LOCAL', or None for ARTIST_ROSTER) country_code: Country code for LOCAL_REP single country deletion delete_all_countries: If True, delete all LOCAL_REP countries Raises: ValueError: If validation fails or PROD writes are blocked """ # CRITICAL: Block all PROD writes for safety in QA environment if IS_PROD: raise ValueError( "PROD writes are BLOCKED for safety. " "Use QA schema for testing. Set target=qa environment variable for QA writes." ) schema_name = queries.SCHEMA_NAME # Build appropriate DELETE query if roster_type == "MAIN": query = f""" DELETE FROM {schema_name}.ARTIST_ROSTER_MAIN_REP WHERE VENDOR_ID = {vendor_id} AND GLOBAL_PARTICIPANT_ID = '{artist_uuid}' AND SUBACCOUNT_ID = {subaccount_id} """ elif roster_type == "LOCAL": if delete_all_countries: query = f""" DELETE FROM {schema_name}.ARTIST_ROSTER_LOCAL_REP WHERE VENDOR_ID = {vendor_id} AND GLOBAL_PARTICIPANT_ID = '{artist_uuid}' AND SUBACCOUNT_ID = {subaccount_id} """ else: if not country_code: raise ValueError( "country_code is required for LOCAL_REP single country deletion" ) query = f""" DELETE FROM {schema_name}.ARTIST_ROSTER_LOCAL_REP WHERE VENDOR_ID = {vendor_id} AND GLOBAL_PARTICIPANT_ID = '{artist_uuid}' AND SUBACCOUNT_ID = {subaccount_id} AND COUNTRY_CODE = '{country_code}' """ else: # roster_type is None (ARTIST_ROSTER) query = f""" DELETE FROM {schema_name}.ARTIST_ROSTER WHERE VENDOR_ID = {vendor_id} AND GLOBAL_PARTICIPANT_ID = '{artist_uuid}' AND SUBACCOUNT_ID = {subaccount_id} """ # Execute delete execute_non_query(session, query) def parse_csv_for_bulk_upload( uploaded_file: Any, roster_type: str ) -> tuple[pd.DataFrame, list[str]]: """ Parse and validate CSV structure for bulk upload. Args: uploaded_file: Streamlit UploadedFile object roster_type: MAIN_REP, LOCAL_REP, or ARTIST_ROSTER Returns: Tuple of (DataFrame, list of parsing errors) """ parse_errors = [] try: # Read CSV with pandas df = pd.read_csv(uploaded_file) # Check row limit if len(df) > 250: parse_errors.append( f"CSV exceeds 250 artist limit (found {len(df)} rows)" ) return pd.DataFrame(), parse_errors if len(df) == 0: parse_errors.append("CSV is empty") return pd.DataFrame(), parse_errors # Validate columns based on roster type required_cols = ["artist_name", "spotify_id"] # Check required columns exist missing_cols = [col for col in required_cols if col not in df.columns] if missing_cols: parse_errors.append( f"Missing required columns: {', '.join(missing_cols)}" ) return pd.DataFrame(), parse_errors # For MAIN_REP: status column is optional if roster_type == "MAIN_REP": if "status" not in df.columns: # Add status column with default ACTIVE df["status"] = "ACTIVE" else: # Validate status values invalid_status = df[~df["status"].isin(["ACTIVE", "INACTIVE"])] if not invalid_status.empty: invalid_rows = ", ".join( map(str, [idx + 2 for idx in invalid_status.index]) ) parse_errors.append( f"Invalid status values on rows: {invalid_rows}. " "Status must be ACTIVE or INACTIVE" ) # Strip whitespace from all string columns for col in df.columns: if df[col].dtype == "object": df[col] = df[col].str.strip() # Check for empty required fields for col in required_cols: empty_rows = df[df[col].isna() | (df[col] == "")] if not empty_rows.empty: empty_row_nums = ", ".join(map(str, [idx + 2 for idx in empty_rows.index])) parse_errors.append(f"Empty {col} on rows: {empty_row_nums}") return df, parse_errors except pd.errors.EmptyDataError: parse_errors.append("CSV file is empty") return pd.DataFrame(), parse_errors except Exception as e: parse_errors.append(f"CSV parsing error: {str(e)}") return pd.DataFrame(), parse_errors def check_existing_roster_entry( session: Session, vendor_id: int, artist_uuid: str, subaccount_id: int, roster_type: str, country_code: str | None = None, status: str = "ACTIVE", ) -> tuple[bool, bool, Any]: """ Check if artist already exists in roster and compare data. Args: session: Snowpark session vendor_id: Vendor ID artist_uuid: Artist UUID subaccount_id: Subaccount ID roster_type: MAIN_REP, LOCAL_REP, or ARTIST_ROSTER country_code: Country code (for LOCAL_REP) status: Status (for MAIN_REP) Returns: Tuple of (exists, data_matches, existing_data) - exists: True if entry exists in roster - data_matches: True if existing data matches new data - existing_data: Dict of existing data or None """ schema_name = queries.SCHEMA_NAME if roster_type == "MAIN_REP": query = f""" SELECT STATUS, IS_ARTIST_TEAM FROM {schema_name}.ARTIST_ROSTER_MAIN_REP WHERE VENDOR_ID = {vendor_id} AND GLOBAL_PARTICIPANT_ID = '{artist_uuid}' AND SUBACCOUNT_ID = {subaccount_id} LIMIT 1 """ df = execute_query(session, query) if df.empty: return False, False, None existing_status = df.iloc[0]["STATUS"] return True, existing_status == status, {"status": existing_status} elif roster_type == "LOCAL_REP": query = f""" SELECT COUNTRY_CODE FROM {schema_name}.ARTIST_ROSTER_LOCAL_REP WHERE VENDOR_ID = {vendor_id} AND GLOBAL_PARTICIPANT_ID = '{artist_uuid}' AND SUBACCOUNT_ID = {subaccount_id} AND COUNTRY_CODE = '{country_code}' LIMIT 1 """ df = execute_query(session, query) if df.empty: return False, False, None # If exists with same country, data matches return True, True, {"country_code": country_code} else: # ARTIST_ROSTER query = f""" SELECT 1 FROM {schema_name}.ARTIST_ROSTER WHERE VENDOR_ID = {vendor_id} AND GLOBAL_PARTICIPANT_ID = '{artist_uuid}' AND SUBACCOUNT_ID = {subaccount_id} LIMIT 1 """ df = execute_query(session, query) if df.empty: return False, False, None return True, True, {} def validate_bulk_artists( session: Session, df: pd.DataFrame, vendor_id: int, roster_type: str, country_code: str | None = None, ) -> tuple[list[AddArtistForm], list[Any], list[Any]]: """ Validate all artists in CSV before insertion. Validation Steps: 1. Parse Spotify IDs and lookup artists 2. Check for duplicate Spotify IDs in CSV 3. Deduplicate identical entries 4. Validate each artist exists in DB 5. Check if artist already exists in roster 6. Collect all errors (no partial processing) Args: session: Snowpark session df: Parsed CSV DataFrame vendor_id: Selected vendor ID roster_type: MAIN_REP, LOCAL_REP, or ARTIST_ROSTER country_code: Country code (required for LOCAL_REP) Returns: Tuple of (valid_forms, validation_errors, already_exists_list) """ from common.types import BulkArtistAlreadyExists, BulkValidationError validation_errors = [] already_exists_list = [] valid_forms = [] seen_spotify_ids: dict[str, str] = {} # spotify_id -> artist_name # Step 1: Parse Spotify IDs df["parsed_spotify_id"] = df["spotify_id"].apply(parse_spotify_id) # Step 2: Check for duplicate Spotify IDs with different names for idx, row in df.iterrows(): row_num = int(idx) + 2 if isinstance(idx, (int, float)) else 2 # CSV row number (header = 1) spotify_id = row["parsed_spotify_id"] artist_name = row["artist_name"] if spotify_id in seen_spotify_ids: existing_name = seen_spotify_ids[spotify_id] if existing_name != artist_name: # Same Spotify ID with different names = ERROR validation_errors.append( BulkValidationError( row_number=row_num, spotify_id=spotify_id, artist_name=artist_name, error_type="DUPLICATE_SPOTIFY_ID", message=( f"Spotify ID '{spotify_id}' appears with different names: " f"'{existing_name}' and '{artist_name}'" ), ) ) else: seen_spotify_ids[spotify_id] = artist_name # If duplicate ID errors found, stop here if validation_errors: return [], validation_errors, [] # Step 3: Deduplicate identical entries (same spotify_id + same name) df["unique_key"] = df.apply( lambda row: ( row["parsed_spotify_id"], row["artist_name"], row.get("status", "ACTIVE") if roster_type == "MAIN_REP" else None, ), axis=1, ) # Keep first occurrence of each unique entry df_deduplicated = df.drop_duplicates(subset=["unique_key"], keep="first") # Update progress try: import streamlit as st if hasattr(st.session_state, "validation_progress_bar"): st.session_state.validation_progress_bar.progress(0.2) st.session_state.validation_progress_text.text( f"🔍 Looking up {len(df_deduplicated)} artists in database..." ) except Exception: pass # Step 4: Batch lookup all artists by Spotify IDs (OPTIMIZED - single query) all_spotify_ids = df_deduplicated["parsed_spotify_id"].tolist() # Build WHERE clause with all Spotify IDs spotify_ids_str = "', '".join(all_spotify_ids) batch_query = f""" SELECT DISTINCT ID AS ARTIST_UUID, NAME AS ARTIST_NAME, SPOTIFY_ID FROM orchard_app_reporting.delphi_prod.GLOBAL_PARTICIPANT WHERE NOT _FIVETRAN_DELETED AND SPOTIFY_ID IN ('{spotify_ids_str}') """ # Execute batch query df_artists = execute_query(session, batch_query) # Create lookup dictionary: spotify_id -> artist data artists_lookup = {} for _, artist_row in df_artists.iterrows(): spotify_id = artist_row["SPOTIFY_ID"] artists_lookup[spotify_id] = { "artist_uuid": artist_row["ARTIST_UUID"], "artist_name": artist_row["ARTIST_NAME"], } # Update progress try: import streamlit as st if hasattr(st.session_state, "validation_progress_bar"): st.session_state.validation_progress_bar.progress(0.4) st.session_state.validation_progress_text.text( f"🔍 Checking existing roster entries..." ) except Exception: pass # Step 5: Batch check existing roster entries (OPTIMIZED) artist_uuids = [data["artist_uuid"] for data in artists_lookup.values()] existing_roster = {} if artist_uuids: artist_uuids_str = "', '".join(artist_uuids) if roster_type == "MAIN_REP": roster_query = f""" SELECT GLOBAL_PARTICIPANT_ID, STATUS, IS_ARTIST_TEAM FROM {queries.SCHEMA_NAME}.ARTIST_ROSTER_MAIN_REP WHERE VENDOR_ID = {vendor_id} AND SUBACCOUNT_ID = 0 AND GLOBAL_PARTICIPANT_ID IN ('{artist_uuids_str}') """ elif roster_type == "LOCAL_REP": roster_query = f""" SELECT GLOBAL_PARTICIPANT_ID, COUNTRY_CODE FROM {queries.SCHEMA_NAME}.ARTIST_ROSTER_LOCAL_REP WHERE VENDOR_ID = {vendor_id} AND SUBACCOUNT_ID = 0 AND GLOBAL_PARTICIPANT_ID IN ('{artist_uuids_str}') """ else: # ARTIST_ROSTER roster_query = f""" SELECT GLOBAL_PARTICIPANT_ID FROM {queries.SCHEMA_NAME}.ARTIST_ROSTER WHERE VENDOR_ID = {vendor_id} AND SUBACCOUNT_ID = 0 AND GLOBAL_PARTICIPANT_ID IN ('{artist_uuids_str}') """ df_existing = execute_query(session, roster_query) # Create lookup dictionary for _, ex_row in df_existing.iterrows(): artist_uuid = ex_row["GLOBAL_PARTICIPANT_ID"] if roster_type == "MAIN_REP": existing_roster[artist_uuid] = { "status": ex_row.get("STATUS"), "is_artist_team": ex_row.get("IS_ARTIST_TEAM"), } elif roster_type == "LOCAL_REP": # For LOCAL_REP, store list of countries if artist_uuid not in existing_roster: existing_roster[artist_uuid] = [] existing_roster[artist_uuid].append(ex_row.get("COUNTRY_CODE")) else: existing_roster[artist_uuid] = True # Step 6: Process each artist with data from batch queries total_artists = len(df_deduplicated) for i, (idx, row) in enumerate(df_deduplicated.iterrows()): row_num = int(idx) + 2 if isinstance(idx, (int, float)) else i + 2 spotify_id = row["parsed_spotify_id"] artist_name = row["artist_name"] # Update progress try: import streamlit as st if hasattr(st.session_state, "validation_progress_bar"): progress = 0.4 + (0.6 * (i + 1) / total_artists) st.session_state.validation_progress_bar.progress(progress) st.session_state.validation_progress_text.text( f"🔍 Processing artist {i + 1}/{total_artists}: {artist_name[:40]}..." ) except Exception: pass # Check if artist found in batch lookup artist_data = artists_lookup.get(spotify_id) if not artist_data: # Artist not found in database validation_errors.append( BulkValidationError( row_number=row_num, spotify_id=spotify_id, artist_name=artist_name, error_type="ARTIST_NOT_FOUND", message=( f"Artist not found in database: " f"'{artist_name}' (Spotify ID: {spotify_id})" ), ) ) else: # Artist found - create form object artist_uuid = artist_data["artist_uuid"] try: form = AddArtistForm( vendor_id=vendor_id, artist_uuid=artist_uuid, subaccount_id=0, # Always 0 roster_type=roster_type, country_code=country_code, status=row.get("status", "ACTIVE"), is_artist_team=False, # Bulk upload doesn't support this ) # Check if artist already exists in roster (using batch data) existing_data = existing_roster.get(artist_uuid) exists = existing_data is not None data_matches = False if exists: if roster_type == "MAIN_REP": # Check if status matches data_matches = existing_data.get("status") == form.status elif roster_type == "LOCAL_REP": # Check if country code exists in list data_matches = country_code in existing_data else: # ARTIST_ROSTER data_matches = True # No additional data to compare if data_matches: # Already exists with same data - SKIP already_exists_list.append( BulkArtistAlreadyExists( row_number=row_num, spotify_id=spotify_id, artist_name=artist_name, roster_type=roster_type, ) ) else: # Already exists with different data - ERROR diff_msg = "" if roster_type == "MAIN_REP" and existing_data: diff_msg = ( f" (existing status: {existing_data.get('status')})" ) validation_errors.append( BulkValidationError( row_number=row_num, spotify_id=spotify_id, artist_name=artist_name, error_type="ALREADY_EXISTS_DIFFERENT_DATA", message=( f"Artist already exists in roster with different data{diff_msg}" ), ) ) else: # Validate form fields field_errors = form.validate_fields() if field_errors: validation_errors.append( BulkValidationError( row_number=row_num, spotify_id=spotify_id, artist_name=artist_name, error_type="VALIDATION_ERROR", message="; ".join(field_errors), ) ) else: valid_forms.append(form) except Exception as e: validation_errors.append( BulkValidationError( row_number=row_num, spotify_id=spotify_id, artist_name=artist_name, error_type="VALIDATION_ERROR", message=str(e), ) ) return valid_forms, validation_errors, already_exists_list def bulk_add_artists_to_roster( session: Session, forms: list[AddArtistForm] ) -> tuple[int, list[str]]: """ Insert multiple artists to roster using individual MERGE statements. Strategy: - Use existing add_artist_to_roster() for each artist - This ensures all validation and MERGE logic is reused - Collects all errors without stopping Args: session: Snowpark session forms: List of validated AddArtistForm objects Returns: Tuple of (success_count, errors) """ success_count = 0 errors = [] for i, form in enumerate(forms): try: # Use existing service function add_artist_to_roster(session, form) success_count += 1 except Exception as e: errors.append(f"Artist {i + 1} ({form.artist_uuid}): {str(e)}") return success_count, errors def check_virtual_participant_name_exists(session: Session, name: str) -> bool: """Check if VP name exists (case-insensitive).""" query = f""" SELECT 1 FROM {queries.SCHEMA_NAME}.VIRTUAL_PARTICIPANT WHERE UPPER(NAME) = UPPER('{name.replace("'", "''")}') LIMIT 1 """ df = execute_query(session, query) return not df.empty def check_spotify_id_in_global_participant(session: Session, spotify_id: str) -> bool: """Check if Spotify ID exists in GLOBAL_PARTICIPANT.""" query = f""" SELECT 1 FROM orchard_app_reporting.delphi_prod.GLOBAL_PARTICIPANT WHERE SPOTIFY_ID = '{spotify_id}' AND NOT _FIVETRAN_DELETED LIMIT 1 """ df = execute_query(session, query) return not df.empty def check_spotify_id_in_virtual_participant(session: Session, spotify_id: str) -> bool: """Check if Spotify ID exists in VIRTUAL_PARTICIPANT.""" query = f""" SELECT 1 FROM {queries.SCHEMA_NAME}.VIRTUAL_PARTICIPANT WHERE SPOTIFY_ID = '{spotify_id}' LIMIT 1 """ df = execute_query(session, query) return not df.empty def validate_virtual_participant_form( session: Session, form: VirtualParticipantForm ) -> list[str]: """Validate form before submission.""" errors = form.validate_fields() # Check vendor exists vendors = search_vendors(session, vendor_id=form.vendor_id) if not vendors: errors.append(f"Vendor ID {form.vendor_id} does not exist") return errors # Check name uniqueness if check_virtual_participant_name_exists(session, form.name): errors.append(f"Virtual participant with name '{form.name}' already exists") # Spotify ID validation if form.type == "PENDING_ARTIST" and form.spotify_id: parsed_id = parse_spotify_id(form.spotify_id) if check_spotify_id_in_global_participant(session, parsed_id): errors.append( f"Artist with Spotify ID '{parsed_id}' already exists. " "Use Add Artist page instead." ) if check_spotify_id_in_virtual_participant(session, parsed_id): errors.append( f"Virtual participant with Spotify ID '{parsed_id}' already exists" ) return errors def add_virtual_participant(session: Session, form: VirtualParticipantForm) -> None: """Create VIRTUAL_PARTICIPANT and linked CUSTOM_LIST.""" if IS_PROD: raise ValueError("PROD writes are BLOCKED for safety") errors = validate_virtual_participant_form(session, form) if errors: raise ValueError("; ".join(errors)) schema_name = queries.SCHEMA_NAME current_user = get_current_user(session) spotify_id_value = parse_spotify_id(form.spotify_id) if form.spotify_id else None try: # Step 1: Create VIRTUAL_PARTICIPANT spotify_param = f"'{spotify_id_value}'" if spotify_id_value else "NULL" vp_query = f""" MERGE INTO {schema_name}.VIRTUAL_PARTICIPANT t USING ( SELECT UUID_STRING() AS ID, '{form.name.replace("'", "''")}' AS NAME, '{form.type}' AS TYPE, {spotify_param} AS SPOTIFY_ID, CURRENT_TIMESTAMP() AS CREATED_AT, CURRENT_TIMESTAMP() AS UPDATED_AT, '{current_user}' AS CREATED_BY, '{current_user}' AS UPDATED_BY ) s ON (UPPER(t.NAME) = UPPER(s.NAME)) WHEN NOT MATCHED THEN INSERT (ID, NAME, TYPE, SPOTIFY_ID, CREATED_AT, UPDATED_AT, CREATED_BY, UPDATED_BY) VALUES (s.ID, s.NAME, s.TYPE, s.SPOTIFY_ID, s.CREATED_AT, s.UPDATED_AT, s.CREATED_BY, s.UPDATED_BY) """ execute_non_query(session, vp_query) # Step 2: Get VP ID get_id_query = f""" SELECT ID FROM {schema_name}.VIRTUAL_PARTICIPANT WHERE UPPER(NAME) = UPPER('{form.name.replace("'", "''")}') LIMIT 1 """ vp_df = execute_query(session, get_id_query) if vp_df.empty: raise ValueError("Failed to retrieve VP ID") vp_id = vp_df.iloc[0]["ID"] # Step 3: Create CUSTOM_LIST cl_query = f""" MERGE INTO {schema_name}.CUSTOM_LIST t USING ( SELECT UUID_STRING() AS ID, '{form.name.replace("'", "''")}' AS NAME, {form.vendor_id} AS VENDOR_ID, {form.subaccount_id} AS SUBACCOUNT_ID, '{vp_id}' AS VIRTUAL_PARTICIPANT_ID ) s ON (t.VIRTUAL_PARTICIPANT_ID = s.VIRTUAL_PARTICIPANT_ID AND t.VENDOR_ID = s.VENDOR_ID) WHEN NOT MATCHED THEN INSERT (ID, NAME, VENDOR_ID, SUBACCOUNT_ID, VIRTUAL_PARTICIPANT_ID) VALUES (s.ID, s.NAME, s.VENDOR_ID, s.SUBACCOUNT_ID, s.VIRTUAL_PARTICIPANT_ID) """ execute_non_query(session, cl_query) except Exception as e: raise ValueError(f"Failed to create virtual participant: {str(e)}") def search_pending_artists( session: Session, name_filter: str | None = None ) -> list[VirtualParticipantRow]: """ Search for pending artists in VIRTUAL_PARTICIPANT. Args: session: Active Snowpark session name_filter: Optional name pattern for ILIKE filter Returns: List of VirtualParticipantRow objects """ schema_name = queries.SCHEMA_NAME if name_filter and name_filter.strip(): escaped_name = name_filter.replace("'", "''") name_clause = f"AND NAME ILIKE '%{escaped_name}%'" else: name_clause = "" query = f""" SELECT ID, NAME, TYPE, SPOTIFY_ID, CREATED_AT, CREATED_BY FROM {schema_name}.VIRTUAL_PARTICIPANT WHERE TYPE = 'PENDING_ARTIST' {name_clause} ORDER BY NAME LIMIT 50 """ df = execute_query(session, query) if df.empty: return [] return [VirtualParticipantRow(**row) for row in df.to_dict("records")] # type: ignore[misc] def check_spotify_id_in_virtual_participant_excluding( session: Session, spotify_id: str, exclude_vp_id: str ) -> bool: """ Check if Spotify ID exists in VIRTUAL_PARTICIPANT excluding a specific record. Needed because the existing check_spotify_id_in_virtual_participant would falsely flag the current record being updated. Args: session: Active Snowpark session spotify_id: Spotify ID to check exclude_vp_id: Virtual participant ID to exclude from check Returns: True if Spotify ID is used by another VP, False otherwise """ query = f""" SELECT 1 AS exists_in_vp FROM {queries.SCHEMA_NAME}.VIRTUAL_PARTICIPANT WHERE SPOTIFY_ID = '{spotify_id}' AND ID != '{exclude_vp_id}' LIMIT 1 """ df = execute_query(session, query) return not df.empty def upgrade_pending_artist( session: Session, form: UpgradePendingArtistForm ) -> None: """ Update SPOTIFY_ID on a PENDING_ARTIST virtual participant. Args: session: Active Snowpark session form: Upgrade form data with VP ID and new Spotify ID Raises: ValueError: If validation fails, PROD writes are blocked, or uniqueness check fails """ if IS_PROD: raise ValueError( "PROD writes are BLOCKED for safety. " "Use QA schema for testing. Set target=qa environment variable for QA writes." ) errors = form.validate_fields() if errors: raise ValueError("; ".join(errors)) # Check Spotify ID uniqueness (excluding current record) if check_spotify_id_in_virtual_participant_excluding( session, form.spotify_id, form.virtual_participant_id ): raise ValueError( f"Spotify ID '{form.spotify_id}' is already used by another virtual participant" ) schema_name = queries.SCHEMA_NAME current_user = get_current_user(session) query = f""" UPDATE {schema_name}.VIRTUAL_PARTICIPANT SET SPOTIFY_ID = '{form.spotify_id}', UPDATED_AT = CURRENT_TIMESTAMP(), UPDATED_BY = '{current_user}' WHERE ID = '{form.virtual_participant_id}' AND TYPE = 'PENDING_ARTIST' """ execute_non_query(session, query)