"""Migrate video track splits from Abo to Perister.""" import logging import os import sqlalchemy from scripts import script_util from scripts.altafonte_migration.utils import ( SF_ABO_SCHEMA, SF_ART_REL_SCHEMA, SF_COLLABORATOR_TABLE_PROD, SF_COLLABORATOR_TABLE_QA, SF_REMAP_TABLE, insert_splits_using_persister, ) def get_all_video_track_splits(environment): """Get all video track splits. Args: environment (str): The operating environment, e.g. QA, PROD. """ snowflake_conn = script_util.snowflake_connection() sf_collaborator_table = ( environment == "prod" and SF_COLLABORATOR_TABLE_PROD or SF_COLLABORATOR_TABLE_QA ) select_sql = f""" SELECT t.id AS identifier, b.percentage / 100 AS split_rate, 2 AS split_type_id, c.id AS collaborator_id, 'NET' AS rate_type, 'Altafonte migration' AS source, r.release_id AS product_id, FROM {SF_ABO_SCHEMA}.beneficiaries b LEFT JOIN {SF_ABO_SCHEMA}.users su ON b.user_id = su.user_id AND su._fivetran_deleted = FALSE LEFT JOIN {SF_ABO_SCHEMA}.users u ON su.account_id = u.user_ID AND u._fivetran_deleted = FALSE LEFT JOIN {SF_ABO_SCHEMA}.videos av ON b.video = av.id AND av._fivetran_deleted = FALSE INNER JOIN {SF_ABO_SCHEMA}.products ap ON ap.id = av.product_id AND ap._fivetran_deleted = FALSE LEFT JOIN {SF_REMAP_TABLE} remap ON remap.original_upc = ap.ean AND remap._fivetran_deleted = FALSE LEFT JOIN {SF_ART_REL_SCHEMA}.releases r ON ( to_varchar(r.upc) = CASE WHEN remap.remapped_upc IS NOT NULL THEN remap.remapped_upc ELSE ap.ean END AND r._fivetran_deleted = FALSE ) LEFT JOIN {SF_ART_REL_SCHEMA}.track t ON r.release_id = t.release_id AND t._fivetran_deleted = FALSE LEFT JOIN {SF_ART_REL_SCHEMA}.vendor v ON v.vendor_id = u.pde_vendor_id AND v._fivetran_deleted = FALSE LEFT JOIN {sf_collaborator_table} c ON c.internal_id = CONCAT('ABO-', TO_CHAR(su.user_ID)) AND c._fivetran_deleted = FALSE WHERE b._fivetran_deleted = FALSE GROUP BY b.id, t.id, av.id, t.isrc, u.pde_vendor_id, su.user_ID, c.id, b.percentage, r.release_id ; """ return ( snowflake_conn.execute( sqlalchemy.text(select_sql), ) .mappings() .all() ) environment = os.environ.get("Environment", "dev") logging.info(f"Got environment: {environment}") logging.info("Migrating video track splits...") logging.info("Getting video track splits...") video_splits = get_all_video_track_splits(environment) logging.info(f"Found {len(video_splits)} video track splits.") logging.info("Inserting video track splits...") insert_splits_using_persister(video_splits)