"""Snowflake connector.""" from typing import Any from typing import Generator from typing import Iterator from pandas import DataFrame from snowflake import connector import config from src.utils.constants import DBT_TABLE_DISTRIBUTION from src.utils.constants import DBT_TABLE_NEIGHBOURING_RIGHTS from src.utils.constants import StatementAttachmentType from src.utils.error_handling import LambdaException ITERATOR_SIZE = 100000 SUBACCOUNT_FILTER = 'AND WFS.SUBACCOUNTID = %(subaccountid)s' SNOWFLAKE_SCHEMA = str(config.SNOWFLAKE_DB_CONFIG.get('schema') or '') def _build_transaction_type_filters(filters: dict | None, params: dict) -> list[str]: """Build transaction type/group filter clauses. Args: filters (dict | None): Filters containing transaction type/group IDs. params (dict): Query parameters dictionary to update. Returns: list[str]: List of filter clause strings. """ filter_clauses: list[str] = [] if not filters: return filter_clauses if filters.get('transaction_type_ids'): filter_clauses.append('AND DTT.TRANSACTIONTYPEID IN (%(transaction_type_ids)s)') params['transaction_type_ids'] = filters['transaction_type_ids'] if filters.get('exclude_transaction_type_ids'): filter_clauses.append('AND DTT.TRANSACTIONTYPEID NOT IN (%(exclude_transaction_type_ids)s)') params['exclude_transaction_type_ids'] = filters['exclude_transaction_type_ids'] return filter_clauses def _query(query: str, params: dict) -> Generator[Any, None, None]: """Extract data from Snowflake. Args: query (str): SQL statement for extracting data from Snowflake. params (dict): Parameters bound to the SQL statement. Returns: Generator: Generator containing dict rows """ with connector.connect(**config.SNOWFLAKE_DB_CONFIG) as connection: with connection.cursor(connector.DictCursor) as cursor: cursor.execute(query, params) yield cursor.rowcount while True: results = cursor.fetchmany(ITERATOR_SIZE) if not results: break for result in results: yield result def _query_one(query: str, params: dict) -> dict | None: """Execute query and fetch a single row from Snowflake. Args: query (str): SQL statement for extracting data from Snowflake. params: Parameters to pass to the query. Returns: dict | None: Single row as dict or None if no results. """ with connector.connect(**config.SNOWFLAKE_DB_CONFIG) as connection: with connection.cursor(connector.DictCursor) as cursor: cursor.execute(query, params) row = cursor.fetchone() return row if row else None def _pandas_query(query: str, params: dict) -> tuple[Iterator[DataFrame], int]: """Run a query, returning a generator with dataframes and the total rows. Args: query (str): SQL statement for extracting data from Snowflake. params (dict): Parameters to pass along with the query. Returns: Generator, int: Generator containing DataFrame rows and total rows count. """ with connector.connect(**config.SNOWFLAKE_DB_CONFIG) as connection: with connection.cursor(connector.DictCursor) as cursor: cursor.execute(query, params) return cursor.fetch_pandas_batches(), cursor.rowcount or 0 def get_distribution_fact_sales( statement_period_id: int, account_id: int, contract_id: int | None, subaccount_id: int | None ) -> tuple[Iterator[DataFrame], int]: """Get distribution fact sales for a statement/account. Args: statement_period_id (int): Statement period to get sales for account_id (int): Account to get sales for contract_id (int): Contract to filter by subaccount_id (int): Subaccount to filter by Returns: Generator: Generator containing dict rows """ filters = '' params = {'account_id': account_id, 'statement_period_id': statement_period_id} if contract_id: filters += 'AND RDD.CONTRACT_ID = %(contract_id)s' params['contract_id'] = contract_id if subaccount_id: filters += 'AND RDD.SUBACCOUNT_ID = %(subaccount_id)s' params['subaccount_id'] = subaccount_id sql = """ WITH TRACK_PERFORMERS AS ( SELECT LABEL_PARTICIPANT_PARTICIPATED_IN_ORCHARD_TRACK.TRACK_ID, LISTAGG(DISTINCT LABEL_PARTICIPANT.NAME, '|') AS PERFORMERS FROM FACTS.{schema}.LABEL_PARTICIPANT_PARTICIPATED_IN_ORCHARD_TRACK JOIN FACTS.{schema}.LABEL_PARTICIPANT ON LABEL_PARTICIPANT_PARTICIPATED_IN_ORCHARD_TRACK.LABEL_PARTICIPANT_ID = LABEL_PARTICIPANT.ID WHERE PARTICIPATED_AS = 'performer' GROUP BY LABEL_PARTICIPANT_PARTICIPATED_IN_ORCHARD_TRACK.TRACK_ID ) SELECT SP1.STATEMENT_PERIOD_NAME, RDD.ACCOUNT_ID, A.ACCOUNT_NAME, RDD.CONTRACT_ID, RDD.TRANSACTION_DATE, DC.COUNTRYNAME, CMM.CUSTOMER_NAME, RDD.SUBDISTRIBUTOR, DR.IMPRINT AS LABEL_IMPRINT, DA.ARTISTNAME, DR.RELEASENAME, DR.VERSION AS RELEASE_VERSION, DR.PRODUCT_CODE, DR.DISPLAY_UPC, DR.MANUFACTURER_UPC, TP.PERFORMERS, DT.TRACKNAME, DT.VERSION AS TRACK_VERSION, RDD.ISRC, RDD.VIDEO_ID, DTT.TRANSACTIONTYPEDESC, RDD.TRANSACTION_SUBTYPE, RDD.UNIT_PRICE_SALE_CURRENCY, RDD.QUANTITY, RDD.GROSS_REVENUE_AFTER_WITHHOLDING_TAX_SALE_CURRENCY, RDD.SALE_CURRENCY_CODE, ER.RATE AS EXCHANGE_RATE, RDD.GROSS_REVENUE_AFTER_WITHHOLDING_TAX_PAYEE_CURRENCY, RDD.ACCOUNT_PAYEE_CURRENCY, RDD.ROYALTY_RATE, RDD.NET_SHARE_PAYEE_CURRENCY, RDD.PHYS_PPD_USD, RDD.MECHANICAL_DEDUCTION_AMOUNT_PAYEE_CURRENCY, RDD.PUBLISHER_ADMIN_FEE_PAYEE_CURRENCY, RDD.ABACUS_SALE_TYPE, SP2.STATEMENT_PERIOD_NAME as ORIGINAL_STATEMENT_PERIOD FROM REVENUE_DISTRO_DBT AS RDD LEFT JOIN FACTS.{schema}.DIM_COUNTRY AS DC ON DC.COUNTRYID = RDD.COUNTRY_ID LEFT JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.CUSTOMER_MASTER_MASTER AS CMM ON CMM.CUSTOMER_MASTER_MASTER_ID = RDD.STORE_ID LEFT JOIN FACTS.{schema}.DIM_RELEASE AS DR ON DR.RELEASEID = RDD.UPC LEFT JOIN ( SELECT DIM_TRACK.* FROM FACTS.{schema}.DIM_TRACK JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.TRACK AS T on T.ID = DIM_TRACK.TRACK_UNIQUE_ID ) AS DT ON DT.UPC = RDD.UPC AND DT.ISRC = RDD.ISRC AND DT.TRACK_ID = RDD.TRACK_ID LEFT JOIN FACTS.{schema}.DIM_ARTIST AS DA ON DA.ARTISTID = DR.ARTISTID LEFT JOIN ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.EXCHANGE_RATE AS ER ON RDD.STATEMENT_PERIOD_ID = ER.STATEMENT_PERIOD_ID AND RDD.SALE_CURRENCY_CODE = ER.FROM_CURRENCY_CODE AND RDD.ACCOUNT_PAYEE_CURRENCY = ER.TO_CURRENCY_CODE LEFT JOIN ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.STATEMENT_PERIOD AS SP1 ON RDD.STATEMENT_PERIOD_ID = SP1.STATEMENT_PERIOD_ID LEFT JOIN ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.ACCOUNT AS A ON RDD.ACCOUNT_ID = A.ACCOUNT_ID LEFT JOIN FACTS.{schema}.DIM_TRANSACTIONTYPE AS DTT ON DTT.TRANSACTIONTYPEABBR = RDD.TRANSACTION_TYPE LEFT JOIN FACTS.{schema}.DIM_LABEL AS DL ON DL.LABELID = RDD.LABEL_ID LEFT JOIN TRACK_PERFORMERS TP ON RDD.TRACK_UNIQUE_ID = TP.TRACK_ID LEFT JOIN ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.STATEMENT_PERIOD AS SP2 ON RDD.ORIGINAL_STATEMENT_PERIOD_ID = SP2.STATEMENT_PERIOD_ID WHERE RDD.STATEMENT_PERIOD_ID = %(statement_period_id)s AND RDD.ACCOUNT_ID = %(account_id)s {filters} """.format(schema=SNOWFLAKE_SCHEMA, filters=filters) return _pandas_query(sql, params) def get_neighbouring_rights_label_fact_sales( statement_period_id: int, account_id: int, contract_id: int ) -> tuple[Generator[Any, None, None], Any]: """Get neighbouring rights label fact sales for a statement/account. Args: statement_period_id (int): Statement period to get sales for account_id (int): Account to get sales for contract_id (int): Contract to get sales for Returns: Generator: Generator containing dict rows """ sql = """ SELECT RDD.ACCOUNT_PAYEE_CURRENCY, RDD.CONTRACT_ID, RDD.GROSS_REVENUE_AFTER_WITHHOLDING_TAX_PAYEE_CURRENCY, RDD.GROSS_REVENUE_PAYEE_CURRENCY, RDD.ISRC, RDD.NET_SHARE_PAYEE_CURRENCY, RDD.ROYALTY_RATE, RDD.START_DATE, RDD.TRANSACTION_DATE, RDD.TRANSACTION_SUBTYPE, RDD.WITHHOLDING_TAX_PAYEE_CURRENCY, DA.ARTISTNAME, DC.COUNTRYNAME, DR.IMPRINT, CMM.CUSTOMER_NAME, DT.TRACKNAME, DT.VERSION, DTT.TRANSACTIONTYPEDESC, SP1.STATEMENT_PERIOD_NAME, RDD.ABACUS_SALE_TYPE, SP2.STATEMENT_PERIOD_NAME as ORIGINAL_STATEMENT_PERIOD FROM REVENUE_DISTRO_DBT AS RDD LEFT JOIN FACTS.{schema}.DIM_COUNTRY AS DC ON DC.COUNTRYID = RDD.COUNTRY_ID LEFT JOIN FACTS.{schema}.DIM_RELEASE AS DR ON DR.RELEASEID = RDD.UPC LEFT JOIN FACTS.{schema}.DIM_ARTIST AS DA ON DA.ARTISTID = DR.ARTISTID LEFT JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.CUSTOMER_MASTER_MASTER AS CMM ON CMM.CUSTOMER_MASTER_MASTER_ID = RDD.STORE_ID LEFT JOIN ( SELECT DIM_TRACK.* FROM FACTS.{schema}.DIM_TRACK JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.TRACK AS T on T.ID = DIM_TRACK.TRACK_UNIQUE_ID ) AS DT ON DT.UPC = RDD.UPC AND DT.ISRC = RDD.ISRC AND DT.TRACK_ID = RDD.TRACK_ID LEFT JOIN FACTS.{schema}.DIM_TRANSACTIONTYPE AS DTT ON DTT.TRANSACTIONTYPEABBR = RDD.TRANSACTION_TYPE LEFT JOIN ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.STATEMENT_PERIOD AS SP1 ON SP1.STATEMENT_PERIOD_ID = RDD.STATEMENT_PERIOD_ID LEFT JOIN ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.STATEMENT_PERIOD AS SP2 ON RDD.ORIGINAL_STATEMENT_PERIOD_ID = SP2.STATEMENT_PERIOD_ID WHERE RDD.STATEMENT_PERIOD_ID = %(statement_period_id)s AND RDD.ACCOUNT_ID = %(account_id)s AND RDD.CONTRACT_ID = %(contract_id)s """.format(schema=SNOWFLAKE_SCHEMA) params = { 'account_id': account_id, 'contract_id': contract_id, 'statement_period_id': statement_period_id, } generator = _query(sql, params) total_rows = next(generator) return generator, total_rows def get_neighbouring_rights_performer_fact_sales( statement_period_id: int, account_id: int, contract_id: int ) -> tuple[Generator[Any, None, None], Any]: """Get neighbouring rights performer fact sales for a statement/account. Args: statement_period_id (int): Statement period to get sales for account_id (int): Account to get sales for contract_id (int): Contract to get sales for Returns: Generator: Generator containing dict rows """ sql = """ SELECT RND.ACCOUNT_PAYEE_CURRENCY, RND.CONTRACT_ID, RND.CONTRIBUTOR_NAME, RND.END_DATE, RND.GROSS_REVENUE_AFTER_WITHHOLDING_TAX_PAYEE_CURRENCY, RND.GROSS_REVENUE_PAYEE_CURRENCY, RND.ISRC, RND.NET_SHARE_PAYEE_CURRENCY, RND.ROYALTY_RATE, RND.SOUND_RECORDING_NAME, RND.START_DATE, RND.TRANSACTION_SUBTYPE_DESCRIPTION, RND.WITHHOLDING_TAX_PAYEE_CURRENCY, CT.CONTRACT_TERM_NAME, CT.TERM_TYPE, SCH.SCHEDULE_NAME, DC.COUNTRYNAME, CMM.CUSTOMER_NAME, DTT.TRANSACTIONTYPEDESC, SP1.STATEMENT_PERIOD_NAME, SR.MAIN_ARTIST, SR.VERSION, RND.ABACUS_SALE_TYPE, SP2.STATEMENT_PERIOD_NAME as ORIGINAL_STATEMENT_PERIOD FROM REVENUE_NR_DBT AS RND LEFT JOIN CONTRACT_TRANSACTION_NR AS CTNR ON RND.CONTRACT_TRANSACTION_NR_CONTRACT_TXN_ID = CTNR.CONTRACT_TXN_ID LEFT JOIN ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.CONTRACT_TERM_CONDITION AS CTC ON CTC.CONTRACT_TERM_CONDITION_ID = RND.CONTRACT_TERM_CONDITION_ID LEFT JOIN ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.CONTRACT_TERM AS CT ON CT.CONTRACT_TERM_ID = CTC.CONTRACT_TERM_ID LEFT JOIN ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.SCHEDULE AS SCH ON SCH.SCHEDULE_ID = CTNR.SCHEDULE_ID LEFT JOIN FACTS.{schema}.DIM_COUNTRY AS DC ON DC.COUNTRYID = RND.COUNTRY_ID LEFT JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.CUSTOMER_MASTER_MASTER AS CMM ON CMM.CUSTOMER_MASTER_MASTER_ID = RND.STORE_ID LEFT JOIN FACTS.{schema}.DIM_TRANSACTIONTYPE AS DTT ON DTT.TRANSACTIONTYPEABBR = RND.TRANSACTION_TYPE LEFT JOIN ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.STATEMENT_PERIOD AS SP1 ON RND.STATEMENT_PERIOD_ID = SP1.STATEMENT_PERIOD_ID LEFT JOIN FACTS.{schema}.PERFORMANCE_NR_SOUND_RECORDING AS SR ON SR.ID = RND.SOUND_RECORDING_ID LEFT JOIN ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.STATEMENT_PERIOD AS SP2 ON RND.ORIGINAL_STATEMENT_PERIOD_ID = SP2.STATEMENT_PERIOD_ID WHERE RND.STATEMENT_PERIOD_ID = %(statement_period_id)s AND RND.ACCOUNT_ID = %(account_id)s AND RND.CONTRACT_ID = %(contract_id)s """.format(schema=SNOWFLAKE_SCHEMA) params = { 'account_id': account_id, 'contract_id': contract_id, 'statement_period_id': statement_period_id, } generator = _query(sql, params) total_rows = next(generator) return generator, total_rows def get_collection_summary_data( statement_period_id: int, account_id: int, contract_id: int, statement_attachment_type: str ) -> tuple[Generator[Any, None, None], Any]: """Get collection summary data for a statement/account. Args: statement_period_id (int): Statement period to get sales for account_id (int): Account to get sales for contract_id (int): Contract to get sales for statement_attachment_type (str): Statement attachment type Returns: Generator: Generator containing dict rows """ match statement_attachment_type: case StatementAttachmentType.COLLECTION_SUMMARY_LABEL: table_name = DBT_TABLE_DISTRIBUTION case StatementAttachmentType.COLLECTION_SUMMARY_PERFORMER: table_name = DBT_TABLE_NEIGHBOURING_RIGHTS case _: raise LambdaException( f'Unhandled statement attachment type: {statement_attachment_type}' ) sql = """ SELECT CUSTOMER_MASTER_MASTER.CUSTOMER_NAME AS COLLECTION_SOCIETY, DIM_COUNTRY.COUNTRYNAME AS COUNTRY, {table_name}.ACCOUNT_PAYEE_CURRENCY AS CURRENCY, SUM({table_name}.GROSS_REVENUE_PAYEE_CURRENCY) AS PRE_WHT_AMOUNT, SUM({table_name}.WITHHOLDING_TAX_PAYEE_CURRENCY) AS WHT_AMOUNT, SUM({table_name}.GROSS_REVENUE_AFTER_WITHHOLDING_TAX_PAYEE_CURRENCY) AS GROSS_AMOUNT, (SELECT {table_name}.ROYALTY_RATE * 100) AS CLIENT_SHARE, SUM({table_name}.NET_SHARE_PAYEE_CURRENCY) AS NET_REVENUE FROM {table_name} JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.CUSTOMER_MASTER_MASTER ON CUSTOMER_MASTER_MASTER.CUSTOMER_MASTER_MASTER_ID = {table_name}.STORE_ID JOIN FACTS.{schema}.DIM_COUNTRY ON DIM_COUNTRY.COUNTRYID = {table_name}.COUNTRY_ID WHERE {table_name}.STATEMENT_PERIOD_ID = %(statement_period_id)s AND {table_name}.ACCOUNT_ID = %(account_id)s AND {table_name}.CONTRACT_ID = %(contract_id)s GROUP BY COUNTRY, COLLECTION_SOCIETY, CURRENCY, CLIENT_SHARE ORDER BY COLLECTION_SOCIETY, COUNTRY ASC """.format(table_name=table_name, schema=SNOWFLAKE_SCHEMA) params = { 'account_id': account_id, 'contract_id': contract_id, 'statement_period_id': statement_period_id, } generator = _query(sql, params) total_rows = next(generator) return generator, total_rows def get_workstation_fact_sales( statement_period_id: int, account_id: int, subaccount_id: int | None ) -> tuple[Generator[dict, None, None], int]: """Get workstation fact sales for a specific statement period and account. Args: statement_period_id (int): Statement period to get sales for account_id (int): Account to get sales for subaccount_id (int): Subaccount to get the sales for Returns: Generator, int: Generator containing dict rows and total row count. """ filters = '' params = {'accountingperiodid': statement_period_id, 'labelid': account_id} if subaccount_id: filters = SUBACCOUNT_FILTER params['subaccountid'] = subaccount_id sql = """ SELECT CONCAT(WFS.ACCOUNTINGYEAR, 'M', WFS.ACCOUNTINGMONTH) AS PERIOD, CONCAT(WFS.ACTIVITYYEAR, 'M', WFS.ACTIVITYMONTH) AS ACTIVITYPERIOD, TRIM(CMM.CUSTOMER_NAME) AS CUSTOMER_NAME, DC.COUNTRYNAME, DR.DISPLAY_UPC, DR.MANUFACTURER_UPC, DR.VENDOR_CATALOG_NUMBER, DR.PRODUCT_CODE, SA.SUBACCOUNT_NAME, IM.IMPRINT, DA.ARTISTNAME, DR.RELEASENAME, COALESCE( CASE WHEN DT.CD = 0 AND DT.TRACK_ID = 0 THEN 'FULL ALBUM' ELSE DT.TRACKNAME END, '' ) AS TRACKNAME, ( SELECT LISTAGG(NAME, '|') WITHIN GROUP(ORDER BY NAME) FROM FACTS.PROD.TRACK_ARTIST TA WHERE TA.TRACK_ID = DT.TRACK_UNIQUE_ID AND TYPE = 'performer' ) AS TRACKARTIST, COALESCE( CASE WHEN DT.CD = 0 AND DT.TRACK_ID = 0 THEN AS_VARCHAR(DR.RELEASEID) ELSE DT.ISRC END, '' ) AS ISRC, COALESCE(DT.CD, 0) AS CD, COALESCE(DT.TRACK_ID, 0) AS TRACK_ID, DTT.TRANSACTIONTYPEABBR, DTT.TRANSACTIONTYPEDESC, WFS.ORIGINAL_PRICE, WFS.DISCOUNT, CASE WHEN WFS.SALES <> 0 THEN (WFS.FX_GROSS / WFS.SALES) ELSE 0 END AS ACTUAL_PRICE, CAST(COALESCE(WFS.SALES, 0) AS INT) AS SALES, WFS.FX_GROSS, WFS.FX_ADJUSTED_GROSS, CASE WHEN WFS.FX_ADJUSTED_GROSS <> 0 THEN (WFS.FX_NET_RECEIPT::FLOAT / WFS.FX_ADJUSTED_GROSS::FLOAT)::DECIMAL(38, 19) ELSE 0 END AS SPLITRATE, WFS.FX_NET_RECEIPT, COALESCE(WFS.FX_RINGTONE_PUBLISHING, 0.0) AS FX_RINGTONE_PUBLISHING, COALESCE(WFS.FX_CLOUD_PUBLISHING, 0.0) AS FX_CLOUD_PUBLISHING, COALESCE(WFS.FX_DPD_PUBLISHING, 0.0) AS FX_DPD_PUBLISHING, COALESCE(WFS.FX_OMS_FEES, 0.0) AS FX_OMS_FEES, DIMC.CURRENCY_CODE FROM WORKSTATION_FACT_SALES_UNIFIED_DBT WFS LEFT JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.CUSTOMER_MASTER_MASTER CMM ON CMM.CUSTOMER_MASTER_MASTER_ID = WFS.STOREID LEFT JOIN FACTS.{schema}.DIM_COUNTRY DC ON DC.COUNTRYID = WFS.COUNTRYID LEFT JOIN FACTS.{schema}.DIM_RELEASE DR ON DR.RELEASEID = WFS.RELEASEID LEFT JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.SUBACCOUNT SA ON SA.SUBACCOUNT_ID = WFS.SUBACCOUNTID LEFT JOIN FACTS.{schema}.DIM_IMPRINT IM ON IM.IMPRINTID = WFS.IMPRINTID LEFT JOIN FACTS.{schema}.DIM_ARTIST DA ON DA.ARTISTID = WFS.ARTISTID LEFT JOIN FACTS.{schema}.DIM_TRACK_CLEAN_MV DT ON DT.TRACKID = WFS.TRACKID LEFT JOIN FACTS.{schema}.DIM_TRANSACTIONTYPE DTT ON DTT.TRANSACTIONTYPEID = WFS.TRANSACTIONTYPEID LEFT JOIN FACTS.{schema}.DIM_CURRENCY DIMC ON WFS.PAYOUT_CURRENCY_ID = DIMC.CURRENCYID WHERE WFS.ACCOUNTINGPERIODID = %(accountingperiodid)s AND WFS.LABELID = %(labelid)s {filters} """.format(schema=SNOWFLAKE_SCHEMA, filters=filters) generator = _query(sql, params) total_rows = next(generator) return generator, total_rows def get_workstation_fact_sales_pandas( statement_period_ids: tuple[int, ...], account_id: int, subaccount_id: int | None, filters: dict | None = None, ) -> tuple[Iterator[DataFrame], int]: """Get workstation fact sales for a specific statement period and account (using Pandas). Args: statement_period_ids (tuple[int]): Statement period ids to get sales for account_id (int): Account to get sales for subaccount_id (int): Subaccount to get the sales for filters (dict | None): Filters to apply to the data Returns: Generator: Generator containing dict rows. """ filter_clauses = [] params = {'accountingperiod_ids': statement_period_ids, 'labelid': account_id} if subaccount_id: filter_clauses.append(SUBACCOUNT_FILTER) params['subaccountid'] = subaccount_id filter_clauses.extend(_build_transaction_type_filters(filters, params)) filters_sql = ' '.join(filter_clauses) sql = """ SELECT CONCAT(WFS.ACCOUNTINGYEAR, 'M', WFS.ACCOUNTINGMONTH) AS PERIOD, CONCAT(WFS.ACTIVITYYEAR, 'M', WFS.ACTIVITYMONTH) AS ACTIVITYPERIOD, TRIM(CMM.CUSTOMER_NAME) AS CUSTOMER_NAME, DC.COUNTRYNAME, DR.DISPLAY_UPC, DR.MANUFACTURER_UPC, DR.VENDOR_CATALOG_NUMBER, DR.PRODUCT_CODE, SA.SUBACCOUNT_NAME, IM.IMPRINT, DA.ARTISTNAME, DR.RELEASENAME, COALESCE( CASE WHEN DT.CD = 0 AND DT.TRACK_ID = 0 AND DT.TRACK_UNIQUE_ID = 0 THEN 'Full Album' ELSE DT.TRACKNAME END, '' ) AS TRACKNAME, ( SELECT LISTAGG(NAME, '|') WITHIN GROUP(ORDER BY NAME) FROM FACTS.PROD.TRACK_ARTIST TA WHERE TA.TRACK_ID = DT.TRACK_UNIQUE_ID AND TYPE = 'performer' ) AS TRACKARTIST, COALESCE( CASE WHEN DT.CD = 0 AND DT.TRACK_ID = 0 AND DT.TRACK_UNIQUE_ID = 0 THEN CAST(DR.RELEASEID as VARCHAR) ELSE DT.ISRC END, '' ) AS ISRC, COALESCE(DT.CD, 0) AS CD, COALESCE(DT.TRACK_ID, 0) AS TRACK_ID, DTT.TRANSACTIONTYPEABBR, DTT.TRANSACTIONTYPEDESC, WFS.ORIGINAL_PRICE, WFS.DISCOUNT, CASE WHEN WFS.SALES <> 0 THEN (WFS.FX_GROSS / WFS.SALES) ELSE 0 END AS ACTUAL_PRICE, CAST(COALESCE(WFS.SALES, 0) AS INT) AS SALES, WFS.FX_GROSS, WFS.FX_ADJUSTED_GROSS, CASE WHEN WFS.FX_ADJUSTED_GROSS <> 0 THEN (WFS.FX_NET_RECEIPT::FLOAT / WFS.FX_ADJUSTED_GROSS::FLOAT)::DECIMAL(38, 19) ELSE 0 END AS SPLITRATE, WFS.FX_NET_RECEIPT, COALESCE(WFS.FX_RINGTONE_PUBLISHING, 0.0) AS FX_RINGTONE_PUBLISHING, COALESCE(WFS.FX_CLOUD_PUBLISHING, 0.0) AS FX_CLOUD_PUBLISHING, COALESCE(WFS.FX_DPD_PUBLISHING, 0.0) AS FX_DPD_PUBLISHING, COALESCE(WFS.FX_OMS_FEES, 0.0) AS FX_OMS_FEES, DIMC.CURRENCY_CODE FROM WORKSTATION_FACT_SALES_UNIFIED_DBT WFS LEFT JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.CUSTOMER_MASTER_MASTER CMM ON CMM.CUSTOMER_MASTER_MASTER_ID = WFS.STOREID LEFT JOIN FACTS.{schema}.DIM_COUNTRY DC ON DC.COUNTRYID = WFS.COUNTRYID LEFT JOIN FACTS.{schema}.DIM_RELEASE DR ON DR.RELEASEID = WFS.RELEASEID LEFT JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.SUBACCOUNT SA ON SA.SUBACCOUNT_ID = WFS.SUBACCOUNTID LEFT JOIN FACTS.{schema}.DIM_IMPRINT IM ON IM.IMPRINTID = WFS.IMPRINTID LEFT JOIN FACTS.{schema}.DIM_ARTIST DA ON DA.ARTISTID = WFS.ARTISTID LEFT JOIN FACTS.{schema}.DIM_TRACK DT ON DT.TRACKID = WFS.TRACKID LEFT JOIN FACTS.{schema}.DIM_TRANSACTIONTYPE DTT ON DTT.TRANSACTIONTYPEID = WFS.TRANSACTIONTYPEID LEFT JOIN FACTS.{schema}.DIM_CURRENCY DIMC ON WFS.PAYOUT_CURRENCY_ID = DIMC.CURRENCYID WHERE WFS.ACCOUNTINGPERIODID IN (%(accountingperiod_ids)s) AND WFS.LABELID = %(labelid)s {filters} """.format(schema=SNOWFLAKE_SCHEMA, filters=filters_sql) return _pandas_query(sql, params) def get_workstation_physical_sales( statement_period_id: int, account_id: int, subaccount_id: int | None ) -> tuple[Generator[dict, None, None], int]: """Get workstation physical sales for a specific statement period and account. Args: statement_period_id (int): Statement period to get sales for. account_id (int): Account to get sales for. subaccount_id (int): Subaccount to get sales for. Returns: Generator, int: Generator containing dict rows and total row count. """ filters = '' params = {'accountingperiodid': statement_period_id, 'labelid': account_id} if subaccount_id: filters = SUBACCOUNT_FILTER params['subaccountid'] = subaccount_id sql = """ SELECT CONCAT(WFS.ACCOUNTINGYEAR, 'M', WFS.ACCOUNTINGMONTH) AS PERIOD, CONCAT(WFS.ACTIVITYYEAR, 'M', WFS.ACTIVITYMONTH) AS ACTIVITYPERIOD, TRIM(CMM.CUSTOMER_NAME) AS CUSTOMER_NAME, DC.COUNTRYNAME, DR.DISPLAY_UPC, DR.VENDOR_CATALOG_NUMBER, DR.PRODUCT_CODE, SA.SUBACCOUNT_NAME, IM.IMPRINT, DA.ARTISTNAME, DR.RELEASENAME, DR.PHYSICAL_PRODUCT_TYPE, DR.PHYSICAL_PRODUCT_FORMAT, DR.DISPLAY_CONFIGURATION, DTT.TRANSACTIONTYPEABBR, DTT.TRANSACTIONTYPEDESC, WFS.ORIGINAL_PRICE, WFS.DISCOUNT, CASE WHEN WFS.SALES <> 0 THEN (WFS.FX_GROSS / WFS.SALES) ELSE 0 END AS ACTUAL_PRICE, CAST(COALESCE(WFS.SALES, 0) AS INT) AS SALES, WFS.FX_GROSS, WFS.FX_ADJUSTED_GROSS, CASE WHEN WFS.FX_ADJUSTED_GROSS <> 0 THEN (WFS.FX_NET_RECEIPT::FLOAT / WFS.FX_ADJUSTED_GROSS::FLOAT)::DECIMAL(38, 19) ELSE 0 END AS SPLITRATE, WFS.FX_NET_RECEIPT, COALESCE(WFS.FX_DPD_PUBLISHING, 0.0) AS FX_DPD_PUBLISHING, COALESCE(WFS.FX_OMS_FEES, 0.0) AS FX_OMS_FEES, DIMC.CURRENCY_CODE FROM WORKSTATION_FACT_SALES_UNIFIED_DBT WFS LEFT JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.CUSTOMER_MASTER_MASTER CMM ON CMM.CUSTOMER_MASTER_MASTER_ID = WFS.STOREID LEFT JOIN FACTS.{schema}.DIM_COUNTRY DC ON DC.COUNTRYID = WFS.COUNTRYID LEFT JOIN FACTS.{schema}.DIM_RELEASE DR ON DR.RELEASEID = WFS.RELEASEID LEFT JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.SUBACCOUNT SA ON SA.SUBACCOUNT_ID = WFS.SUBACCOUNTID LEFT JOIN FACTS.{schema}.DIM_IMPRINT IM ON IM.IMPRINTID = WFS.IMPRINTID LEFT JOIN FACTS.{schema}.DIM_ARTIST DA ON DA.ARTISTID = WFS.ARTISTID LEFT JOIN FACTS.{schema}.DIM_TRANSACTIONTYPE DTT ON DTT.TRANSACTIONTYPEID = WFS.TRANSACTIONTYPEID LEFT JOIN FACTS.{schema}.DIM_CURRENCY DIMC ON WFS.PAYOUT_CURRENCY_ID = DIMC.CURRENCYID WHERE WFS.ACCOUNTINGPERIODID = %(accountingperiodid)s AND WFS.LABELID = %(labelid)s AND DTT.TRANSACTIONTYPEABBR IN ( 'PS', 'PH', 'OP', 'LP', 'MP', 'TP', 'RE', 'OR', 'LR', 'MR', 'TR' ) {filters} """.format(schema=SNOWFLAKE_SCHEMA, filters=filters) generator = _query(sql, params) total_rows = next(generator) return generator, total_rows def get_workstation_physical_sales_pandas( statement_period_ids: tuple[int, ...], account_id: int, subaccount_id: int | None, filters: dict | None = None, ) -> tuple[Iterator[DataFrame], int]: """Get workstation physical sales for a specific statement period and account (using Pandas). Args: statement_period_ids (tuple[int]): Statement period ids to get sales for. account_id (int): Account to get sales for. subaccount_id (int): Subaccount to get sales for. filters (dict | None): Filters to apply to the data. Returns: Generator, int: Generator containing dict rows and total row count. """ filter_clauses = [] params = {'accountingperiod_ids': statement_period_ids, 'labelid': account_id} if subaccount_id: filter_clauses.append(SUBACCOUNT_FILTER) params['subaccountid'] = subaccount_id filter_clauses.extend(_build_transaction_type_filters(filters, params)) filters_sql = ' '.join(filter_clauses) sql = """ SELECT CONCAT(WFS.ACCOUNTINGYEAR, 'M', WFS.ACCOUNTINGMONTH) AS PERIOD, CONCAT(WFS.ACTIVITYYEAR, 'M', WFS.ACTIVITYMONTH) AS ACTIVITYPERIOD, TRIM(CMM.CUSTOMER_NAME) AS CUSTOMER_NAME, DC.COUNTRYNAME, DR.DISPLAY_UPC, DR.VENDOR_CATALOG_NUMBER, DR.PRODUCT_CODE, SA.SUBACCOUNT_NAME, IM.IMPRINT, DA.ARTISTNAME, DR.RELEASENAME, DR.PHYSICAL_PRODUCT_TYPE, DR.PHYSICAL_PRODUCT_FORMAT, DR.DISPLAY_CONFIGURATION, DTT.TRANSACTIONTYPEABBR, DTT.TRANSACTIONTYPEDESC, WFS.ORIGINAL_PRICE, WFS.DISCOUNT, CASE WHEN WFS.SALES <> 0 THEN (WFS.FX_GROSS / WFS.SALES) ELSE 0 END AS ACTUAL_PRICE, CAST(COALESCE(WFS.SALES, 0) AS INT) AS SALES, WFS.FX_GROSS, WFS.FX_ADJUSTED_GROSS, CASE WHEN WFS.FX_ADJUSTED_GROSS <> 0 THEN (WFS.FX_NET_RECEIPT::FLOAT / WFS.FX_ADJUSTED_GROSS::FLOAT)::DECIMAL(38, 19) ELSE 0 END AS SPLITRATE, WFS.FX_NET_RECEIPT, COALESCE(WFS.FX_DPD_PUBLISHING, 0.0) AS FX_DPD_PUBLISHING, COALESCE(WFS.FX_OMS_FEES, 0.0) AS FX_OMS_FEES, DIMC.CURRENCY_CODE FROM WORKSTATION_FACT_SALES_UNIFIED_DBT WFS LEFT JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.CUSTOMER_MASTER_MASTER CMM ON CMM.CUSTOMER_MASTER_MASTER_ID = WFS.STOREID LEFT JOIN FACTS.{schema}.DIM_COUNTRY DC ON DC.COUNTRYID = WFS.COUNTRYID LEFT JOIN FACTS.{schema}.DIM_RELEASE DR ON DR.RELEASEID = WFS.RELEASEID LEFT JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.SUBACCOUNT SA ON SA.SUBACCOUNT_ID = WFS.SUBACCOUNTID LEFT JOIN FACTS.{schema}.DIM_IMPRINT IM ON IM.IMPRINTID = WFS.IMPRINTID LEFT JOIN FACTS.{schema}.DIM_ARTIST DA ON DA.ARTISTID = WFS.ARTISTID LEFT JOIN FACTS.{schema}.DIM_TRANSACTIONTYPE DTT ON DTT.TRANSACTIONTYPEID = WFS.TRANSACTIONTYPEID LEFT JOIN FACTS.{schema}.DIM_CURRENCY DIMC ON WFS.PAYOUT_CURRENCY_ID = DIMC.CURRENCYID WHERE WFS.ACCOUNTINGPERIODID IN (%(accountingperiod_ids)s) AND WFS.LABELID = %(labelid)s AND DTT.TRANSACTIONTYPEABBR IN ( 'PS', 'PH', 'OP', 'LP', 'MP', 'TP', 'RE', 'OR', 'LR', 'MR', 'TR' ) {filters} """.format(schema=SNOWFLAKE_SCHEMA, filters=filters_sql) return _pandas_query(sql, params) def get_subaccount(subaccount_id: int) -> dict | None: """Get a subaccount by ID.""" sql = """ SELECT SUBACCOUNTNAME, COMMISSIONOVERRIDE, SUBACCOUNT_SPLIT_TYPE FROM FACTS.{schema}.DIM_SUBACCOUNT WHERE SUBACCOUNTID = %(subaccount_id)s """.format(schema=SNOWFLAKE_SCHEMA) params = { 'subaccount_id': subaccount_id, } return _query_one(sql, params) def get_transaction_types_names(transaction_type_ids: list[int]) -> list[str] | None: """Get transaction type names by transaction type IDs. Args: transaction_type_ids (list[int]): List of transaction type IDs Returns: list[str]: List of transaction type names """ sql = """ SELECT TRANSACTIONTYPEDESC FROM FACTS.{schema}.DIM_TRANSACTIONTYPE DTT WHERE TRANSACTIONTYPEID IN (%(transaction_type_ids)s) """.format(schema=SNOWFLAKE_SCHEMA) gen = _query(sql, {'transaction_type_ids': transaction_type_ids}) next(gen) transaction_type_names = [] for row in gen: name = row.get('TRANSACTIONTYPEDESC') if name: transaction_type_names.append(name) if not transaction_type_names: return None return transaction_type_names def get_statement_periods_name(statement_period_id: int) -> dict | None: """Get a subaccount by ID.""" sql = """ SELECT STATEMENT_PERIOD_NAME FROM ORCHARD_APP_REPORTING_V2.{schema}_ROYALTY_ACCOUNTING_ROYALTY_ACCOUNTING.\ STATEMENT_PERIOD WHERE STATEMENT_PERIOD_ID = %(statement_period_id)s """.format(schema=SNOWFLAKE_SCHEMA) params = { 'statement_period_id': statement_period_id, } return _query_one(sql, params) def is_distributor(account_id: int) -> dict | None: """Check if account is a distributor.""" sql = """ SELECT IS_DISTRIBUTOR FROM FACTS.{schema}.VENDOR WHERE VENDOR_ID = %(vendor_id)s """.format(schema=SNOWFLAKE_SCHEMA) params = { 'vendor_id': account_id, } return _query_one(sql, params)