from datetime import datetime, date from snowflake.snowpark.functions import col, concat_ws, lit, coalesce import pandas as pd from typing import Tuple, Optional, List from snowflake.snowpark import Session, FileOperation import streamlit as st from streamlit.runtime.state import SessionStateProxy from streamlit.runtime.uploaded_file_manager import UploadedFile from common.constants_general import ( fan_export_selection_columns, fan_export_column_mappings, advanced_filter_operators_dict, ) from common.constants_general import ( DATABASE_NAME, SCHEMA_NAME, FORMS, TLAS, TLS, LABELS, TERRITORIES, ARTISTS, MAILING_LISTS, APPROVAL_STAGE, FANS, ) @st.cache_data(show_spinner="preparing application...") def prepare_streamlit_app(_snowpark_session: Session) -> Tuple[pd.DataFrame, pd.DataFrame]: """ Function to query data used for populating Streamlit app filter fields :param _snowpark_session: Active Snowpark session :return: Pandas dataframe containing related data """ FORMS_DF = ( _snowpark_session.table(f"{DATABASE_NAME}.{SCHEMA_NAME}.{FORMS}").select( col("id").alias("form_id"), col("form_name_c"), col("tla_id_c"), col("migration_id_c").alias("form_migration_id"), col("last_modified_date").alias("form_last_modified_date"), ) # .filter(col("active_c") == True) # noqa ) MAILING_LIST_DF = _snowpark_session.table( f"{DATABASE_NAME}.{SCHEMA_NAME}.{MAILING_LISTS}" ).select( col("id").alias("mailing_list_id"), col("mailing_list_name_c"), col("tla_id_c"), ) TLAS_DF = _snowpark_session.table(f"{DATABASE_NAME}.{SCHEMA_NAME}.{TLAS}").select( col("id").alias("tla_id"), col("tl_id_c"), col("label_id_c"), col("artist_id_c"), col("territory_id_c"), ) TLS_DF = _snowpark_session.table(f"{DATABASE_NAME}.{SCHEMA_NAME}.{TLS}") LABELS_DF = _snowpark_session.table(f"{DATABASE_NAME}.{SCHEMA_NAME}.{LABELS}").select( col("name").alias("label_name"), col("id").alias("label_id") ) TERRITORIES_DF = _snowpark_session.table(f"{DATABASE_NAME}.{SCHEMA_NAME}.{TERRITORIES}") ARTISTS_DF = _snowpark_session.table(f"{DATABASE_NAME}.{SCHEMA_NAME}.{ARTISTS}").select( col("id").alias("artist_id"), col("artist_c") ) forms_df: pd.DataFrame = ( FORMS_DF.join(TLAS_DF, FORMS_DF.tla_id_c == TLAS_DF.tla_id) .join(TLS_DF, TLAS_DF.tl_id_c == TLS_DF.id) .join(LABELS_DF, TLAS_DF.label_id_c == LABELS_DF.label_id) .join(TERRITORIES_DF, TLS_DF.territory_id_c == TERRITORIES_DF.id) .join(ARTISTS_DF, TLAS_DF.artist_id_c == ARTISTS_DF.artist_id) .select( "form_id", "tla_id", "territory_c", "label_name", "artist_c", "form_name_c", "form_migration_id", "form_last_modified_date", ) .with_column( "migration_id_form_name", concat_ws( lit(" - "), coalesce(col("form_migration_id").cast("string"), lit("")), col("form_name_c"), ), ) .to_pandas() ) mailing_lists_df: pd.DataFrame = ( MAILING_LIST_DF.join(TLAS_DF, MAILING_LIST_DF.tla_id_c == TLAS_DF.tla_id) .join(LABELS_DF, TLAS_DF.label_id_c == LABELS_DF.label_id) .join(TERRITORIES_DF, TLAS_DF.territory_id_c == TERRITORIES_DF.id) .join(ARTISTS_DF, TLAS_DF.artist_id_c == ARTISTS_DF.artist_id) .select( "mailing_list_id", "tla_id", "territory_c", "label_name", "artist_c", "mailing_list_name_c", ) .to_pandas() ) return forms_df, mailing_lists_df @st.dialog("Confirmation") def show_dialog() -> None: """ Function to display Streamlit pop up """ st.write("All successfully finished!") if st.button("OK"): st.session_state["dialog_shown"] = False st.rerun() def execute_audit_flow( executed_query: str, user_filters: str, streamlit_user: str, id: str, df: pd.DataFrame, session: Session, pdf: UploadedFile, page: str, ) -> None: """ Function to store audit data and upload approval :param executed_query: SQL representing query that generated data for rexport :param user_filters: Filters used to create an export :param streamlit_user: Streamlit app user running the export :param id: Unique identifier of export :param df: Pandas dataframe containing fan related export data for audit :param session: Active Snowpark session :param pdf: PDF file containing approval :param page: What Streamlit page called this function :return: Pandas dataframe containing related data """ main_audit_info = [(streamlit_user, id, datetime.now(), executed_query, user_filters)] main_audit_info_df = pd.DataFrame( main_audit_info, columns=["username", "export_id", "export_dt", "query", "selected_filters"] ) session.create_dataframe(main_audit_info_df).write.save_as_table( "fan_data_export_main_audit", mode="append" ) session.create_dataframe(df).write.save_as_table("fan_data_export_values_audit", mode="append") # https://daniel-mannino.medium.com/streamlit-to-upload-files-to-snowflake-stages-e98d1d9e3416 with st.spinner("Uploading approval...."): FileOperation(session).put_stream( input_stream=pdf, auto_compress=False, stage_location="@" + APPROVAL_STAGE + "/" + str(id) + "/" + pdf.name, ) if page == "main page": st.session_state.upload_completed = True st.session_state.file_upload_completed = True st.session_state["dialog_shown"] = True if st.session_state["dialog_shown"]: show_dialog() else: st.session_state.fan_export_upload_completed = True def extract_fan_data( input_df: pd.DataFrame, session: Session, streamlit_session: SessionStateProxy, ) -> tuple[pd.DataFrame, str, bool]: """ Function to query fan related data :param input_df: One column Pandas dataframe containing fan ids :param session: Snowpark session to use for running query in database :param streamlit_session: Streamlit session keys and values :return: Pandas dataframe containing fetched results and used query and indicator for email column usage """ true_keys = [key for key, value in streamlit_session.items() if value] columns_to_use = [ fan_export_selection_columns[str(key)] for key in true_keys if str(key) in fan_export_selection_columns ] columns_to_use_ordered = sorted( columns_to_use, key=lambda x: list(fan_export_selection_columns.values()).index(x) ) if "EMAIL_C" not in columns_to_use_ordered: columns_to_use_ordered.append("EMAIL_C") drop_cols = ["EXPORT_DATA", "ID", "EMAIL"] exclude_email_col = True else: drop_cols = ["EXPORT_DATA", "ID"] exclude_email_col = False if "ADDRESS_1_C" in columns_to_use_ordered: index = columns_to_use_ordered.index("ADDRESS_1_C") columns_to_use_ordered.insert(index + 1, "ADDRESS_2_C") FANS_DF = session.table(f"{DATABASE_NAME}.{SCHEMA_NAME}.{FANS}").select( ["ID"] + columns_to_use_ordered ) FAN_IDS = session.createDataFrame(input_df) JOINED_DF = FAN_IDS.join(FANS_DF, FAN_IDS.fan_id == FANS_DF.id) EXECUTED_DF: pd.DataFrame = JOINED_DF.to_pandas() EXECUTED_DF["EXPORT_DATA"] = EXECUTED_DF.apply(lambda row: row.to_dict(), axis=1) query = JOINED_DF.queries["queries"][0] # print(query) columns_renaming = { k: v for k, v in fan_export_column_mappings.items() if k in EXECUTED_DF.columns } EXECUTED_DF.rename(columns=columns_renaming, inplace=True) st.write( f"Total of {EXECUTED_DF.shape[0]} records. Following columns will be exported: {EXECUTED_DF.drop(columns=drop_cols).columns.tolist()}." ) return EXECUTED_DF, query, exclude_email_col def extract_contest_winners_data( form_id_value: str, contest_time_frame: tuple[date, date], nr_of_winners: int, bonus_entry_selected: bool, session_2: Session, session_state: SessionStateProxy, custom_field_params: List[str], ) -> tuple[pd.DataFrame, str, bool]: """ Function to return contest winners :param form_id_value: One column Pandas dataframe containing fan ids :param contest_time_frame: Tuple containing contest timeframe values(inclusive) :param nr_of_winners: Nr of winners to return :param bonus_entry_selected: Whether bonus entry is selected :param session_2: Snowpark session object for running query in database :param session_state: Streamlit session keys and values used to determine what fan data to extract :param custom_field_params: Cuustom field parameters used in query :return: Pandas dataframe containing fetched results and used query and indicator for email column usage """ true_keys = [key for key, value in session_state.items() if value] columns_to_use = [ fan_export_selection_columns[str(key)] for key in true_keys if str(key) in fan_export_selection_columns ] columns_to_use_ordered = sorted( columns_to_use, key=lambda x: list(fan_export_selection_columns.values()).index(x) ) if "EMAIL_C" not in columns_to_use_ordered: columns_to_use_ordered.append("EMAIL_C") drop_cols = ["EXPORT_DATA", "ID", "EMAIL"] exclude_email_col = True else: drop_cols = ["EXPORT_DATA", "ID"] exclude_email_col = False if "ADDRESS_1_C" in columns_to_use_ordered: index = columns_to_use_ordered.index("ADDRESS_1_C") columns_to_use_ordered.insert(index + 1, "ADDRESS_2_C") filter_columns = "" # first custom field if custom_field_params[0] != "" and custom_field_params[2] != "": if custom_field_params[1] == "Contains": filter_columns += f"AND CONTAINS(lower(FR.{custom_field_params[0]}),lower('{custom_field_params[2]}')) " elif custom_field_params[1] == "Starts with": filter_columns += f"AND STARTSWITH(lower(FR.{custom_field_params[0]}),lower('{custom_field_params[2]}')) " else: operator_value = advanced_filter_operators_dict[custom_field_params[1]] filter_columns += f"AND lower(FR.{custom_field_params[0]}){operator_value}lower('{custom_field_params[2]}') " # second custom field if custom_field_params[3] != "" and custom_field_params[5] != "": if custom_field_params[4] == "Contains": filter_columns += f"AND CONTAINS(lower(FR.{custom_field_params[3]}),lower('{custom_field_params[5]}')) " elif custom_field_params[1] == "Starts with": filter_columns += f"AND STARTSWITH(lower(FR.{custom_field_params[3]}),lower('{custom_field_params[5]}')) " else: operator_value = advanced_filter_operators_dict[custom_field_params[4]] filter_columns += f"AND lower(FR.{custom_field_params[3]}){operator_value}lower('{custom_field_params[5]}') " form_id_values = [a.strip() for a in form_id_value.split(",")] quoted_form_id_values = "', '".join(form_id_values) if bonus_entry_selected: query = f""" WITH FANS_IN_PRIZE_DRAW AS ( SELECT DISTINCT FAN_ID_C FROM {DATABASE_NAME}.{SCHEMA_NAME}.FORM_RESPONSE_C FR INNER JOIN {DATABASE_NAME}.{SCHEMA_NAME}.FORM_C F ON FR.FORM_ID_C = F.ID WHERE ( F.ID IN ('{quoted_form_id_values}') OR F.MIGRATION_ID_C::TEXT IN ('{quoted_form_id_values}') ) {filter_columns} -- can be used for testing: a0QTy00000EwiGkMAJ / 568738 -- can be used for testing: a0QTy00000DwisPMAR/ 554916 ), SUBSCRIBED_FANS AS ( SELECT DISTINCT SUBSCRIPTION_C.FAN_ID_C FROM {DATABASE_NAME}.{SCHEMA_NAME}.FORM_C INNER JOIN {DATABASE_NAME}.{SCHEMA_NAME}.MAILING_LIST_C ON FORM_C.TLA_ID_C = MAILING_LIST_C.TLA_ID_C -- NB! joining all mailing lists of a given artist in a given territory-label INNER JOIN {DATABASE_NAME}.{SCHEMA_NAME}.SUBSCRIPTION_C ON MAILING_LIST_C.ID = SUBSCRIPTION_C.MAILING_LIST_ID_C WHERE ( FORM_C.ID IN ('{quoted_form_id_values}') OR FORM_C.MIGRATION_ID_C::TEXT IN ('{quoted_form_id_values}') ) AND SUBSCRIPTION_C.ACTIVE_C -- OR SUBSCRIPTION_C.LAST_MODIFIED_DATE <= '2025-03-30') ), UNSUBSCRIBED_BUT_ELIGIBLE_FANS AS ( SELECT DISTINCT SUBSCRIPTION_C.FAN_ID_C FROM {DATABASE_NAME}.{SCHEMA_NAME}.FORM_C INNER JOIN {DATABASE_NAME}.{SCHEMA_NAME}.MAILING_LIST_C ON FORM_C.TLA_ID_C = MAILING_LIST_C.TLA_ID_C -- NB! joining all mailing lists of a given artist in a given territory-label INNER JOIN {DATABASE_NAME}.{SCHEMA_NAME}.SUBSCRIPTION_C ON MAILING_LIST_C.ID = SUBSCRIPTION_C.MAILING_LIST_ID_C WHERE ( FORM_C.ID IN ('{quoted_form_id_values}') OR FORM_C.MIGRATION_ID_C::TEXT IN ('{quoted_form_id_values}') ) AND SUBSCRIPTION_C.ACTIVE_C = FALSE AND SUBSCRIPTION_C.LAST_MODIFIED_DATE >= '{contest_time_frame[0]}' -- '2025-03-14' AND SUBSCRIPTION_C.LAST_MODIFIED_DATE <= '{contest_time_frame[1]}' --'2025-03-30' AND SUBSCRIPTION_C.FAN_ID_C IS NOT NULL ), FANS_WITH_SUBSCRIPTION AS ( SELECT * FROM SUBSCRIBED_FANS UNION SELECT * FROM UNSUBSCRIBED_BUT_ELIGIBLE_FANS ), FANS AS ( SELECT FANS_IN_PRIZE_DRAW.FAN_ID_C, TRUE AS PRIZE_DRAW_PARTICIPATION, IFF(FANS_WITH_SUBSCRIPTION.FAN_ID_C IS NULL, FALSE, TRUE) AS ML_SUBSCRIPTION FROM FANS_IN_PRIZE_DRAW LEFT JOIN FANS_WITH_SUBSCRIPTION ON FANS_IN_PRIZE_DRAW.FAN_ID_C = FANS_WITH_SUBSCRIPTION.FAN_ID_C ), FAN_TICKETS AS ( SELECT FAN_ID_C FROM FANS WHERE PRIZE_DRAW_PARTICIPATION = TRUE UNION ALL SELECT FAN_ID_C FROM FANS WHERE ML_SUBSCRIPTION = TRUE ORDER BY FAN_ID_C ) SELECT FAN_ID_C AS FAN_ID FROM FAN_TICKETS ORDER BY RANDOM() LIMIT {nr_of_winners} """ else: query = f""" WITH FANS_IN_PRIZE_DRAW AS ( SELECT DISTINCT FAN_ID_C FROM {DATABASE_NAME}.{SCHEMA_NAME}.FORM_RESPONSE_C FR INNER JOIN {DATABASE_NAME}.{SCHEMA_NAME}.FORM_C F ON FR.FORM_ID_C = F.ID WHERE ( F.ID IN ('{quoted_form_id_values}') OR F.MIGRATION_ID_C::TEXT IN ('{quoted_form_id_values}') ) {filter_columns} ) SELECT FAN_ID_C AS FAN_ID FROM FANS_IN_PRIZE_DRAW ORDER BY RANDOM() LIMIT {nr_of_winners} """ selected_fans = session_2.sql(query).to_pandas() FAN_IDS = session_2.createDataFrame(selected_fans) FANS_DF = session_2.table(f"{DATABASE_NAME}.{SCHEMA_NAME}.{FANS}").select( ["ID"] + columns_to_use_ordered ) JOINED_DF = FAN_IDS.join(FANS_DF, FAN_IDS.fan_id == FANS_DF.id) EXECUTED_DF: pd.DataFrame = JOINED_DF.to_pandas() EXECUTED_DF["EXPORT_DATA"] = EXECUTED_DF.apply(lambda row: row.to_dict(), axis=1) columns_renaming = { k: v for k, v in fan_export_column_mappings.items() if k in EXECUTED_DF.columns } EXECUTED_DF.rename(columns=columns_renaming, inplace=True) st.write( f"Total of {EXECUTED_DF.shape[0]} records. Following columns will be exported: {EXECUTED_DF.drop(columns=drop_cols).columns.tolist()}." ) return EXECUTED_DF, query, exclude_email_col def initiate_data_export( fan_id_source_index: int, snow_session: Session, file_upload: Optional[UploadedFile] = None, csv_header: bool = True, form_id_value: Optional[str] = None, contest_time_frame: Optional[tuple[date, date]] = None, nr_of_winners: Optional[int] = None, bonus_entry_selected: bool = True, custom_field_params: Optional[List[str]] = None, ) -> tuple[pd.DataFrame, pd.DataFrame, str, bool, str]: """ Function to trigger data import based on the selected purpose :param fan_id_source_index: 0 for fans inside csv, 1 for contest winner selection :param snow_session: Snowpark session object for a running query in database :param file_upload: Streamlit fileupload object containing file :param csv_header: Whether CSV file contains header or not :param form_id_value: Form id used to find potential winners :param contest_time_frame: Timeframe used to narrow down data :param nr_of_winners: Nr of winners to select :param bonus_entry_selected: Whether to increase fans chances based on additional data :param custom_field_params: List of custom field parameters used in query :return: Pandas dataframes containing prepared fan data, used query, column usage indicator and export id """ if fan_id_source_index == 1: fan_data_export_id = "FAN_DATA-" + datetime.now().strftime("%Y-%m-%d_%H.%M.%S") if file_upload is not None: if csv_header: fan_id_df = pd.read_csv(file_upload, delimiter=",") else: fan_id_df = pd.read_csv(file_upload, delimiter=",", header=None, names=["FAN_ID"]) else: # Handle the case when file_upload is None raise ValueError("file_upload cannot be None for upload with fans in CSV!") fan_id_df = fan_id_df.set_axis(["FAN_ID"], axis=1) returned_fan_df, query, ignore_email_col = extract_fan_data( fan_id_df, snow_session, st.session_state ) else: fan_data_export_id = "CONTEST_WINNERS-" + datetime.now().strftime("%Y-%m-%d_%H.%M.%S") if ( form_id_value is not None and contest_time_frame is not None and nr_of_winners is not None ): if custom_field_params is None: custom_params = [] else: custom_params = custom_field_params returned_fan_df, query, ignore_email_col = extract_contest_winners_data( form_id_value, contest_time_frame, nr_of_winners, bonus_entry_selected, snow_session, st.session_state, custom_params, ) else: # Handle the case when form_id_value is None raise ValueError("Cannot extract fan data without form id!") st.markdown( f"New export created: {fan_data_export_id}", help="Export id follows this format: EXPORT TYPE-YYYY-MM-DD_HH.MM.SS", ) returned_fan_df["EXPORT_ID"] = fan_data_export_id main_fan_export_df = returned_fan_df.drop(columns=["EXPORT_DATA", "ID", "EXPORT_ID"]) returned_fan_df = returned_fan_df[["EXPORT_ID", "ID", "EXPORT_DATA"]] return main_fan_export_df, returned_fan_df, query, ignore_email_col, fan_data_export_id