"""Unit tests for youtube_asset_conflict specific Snowflake SQL executor.""" from unittest.mock import MagicMock import moto from feed_ingestion.flows.youtube_asset_conflict import config from feed_ingestion.flows.youtube_asset_conflict import snowflake_executor from tests.flows.youtube_asset_conflict.fixtures import ows_territories db_schema_table = '{db}.{schema}.{table}' def _find_execute_query_in_mock_calls(mock_calls): """Return query ran in the execute call.""" for name, args, _ in mock_calls: if 'execute' in name: return args def test_create_staging_raw_temp_table(mocker, sf_config_mock): """Test create_staging_raw_temp_table function in snowflake_executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) staging_raw_temp = config.snowflake_table_names.get( 'staging_raw_temp').format(account='account', datestamp='20170721') snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).create_staging_raw_temp_table(staging_raw_temp) execute_query, _ = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert staging_raw_temp in execute_query else: assert False @moto.mock_aws def test_load_staging_raw_temp_table(mocker, sf_config_mock): """Test load_staging_raw_temp_table function in snowflake_executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) key_dir = 'path' staging_raw_temp = config.snowflake_table_names.get( 'staging_raw_temp').format(account='account', datestamp='20170721') snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).load_staging_raw_temp_table( staging_raw_temp, key_dir) execute_query, query_params = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert staging_raw_temp in execute_query assert key_dir in query_params.get('s3_path') else: assert False def test_truncate_table(mocker, sf_config_mock): """Test truncate_table function in snowflake_executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) staging_raw = config.snowflake_table_names.get('staging_raw') snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).truncate_table(staging_raw) execute_query, _ = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert staging_raw in execute_query else: assert False def test_insert_into_staging_raw_table(mocker, sf_config_mock): """Test insert_into_staging_raw_table in snowflake_executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) staging_raw = config.snowflake_table_names.get('staging_raw') staging_raw_temp = config.snowflake_table_names.get( 'staging_raw_temp').format(account='account', datestamp='20170721') file_name = 'file_name' file_size = 'file_size' download_date = 'download_date' content_owner = 'content_owner' snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).insert_into_staging_raw_table( staging_raw_temp, staging_raw, file_name, file_size, download_date, content_owner) execute_query, query_params = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert staging_raw in execute_query assert staging_raw_temp in execute_query assert file_name in query_params.get(file_name) assert file_size in query_params.get(file_size) assert download_date in query_params.get(download_date) assert content_owner in query_params.get(content_owner) else: assert False def test_create_territories_temp_table(mocker, sf_config_mock): """Test create_territories_temp_table in snowflake executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) executor = snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock) tmp_territories_table_name = config.snowflake_table_names[ 'territories_temp'] executor.create_territories_temp_table(tmp_territories_table_name) execute_query, _ = _find_execute_query_in_mock_calls( mock_connect.mock_calls) expected_table_name = '{db}.{schema}.{table}'.format( db=sf_config_mock['db'], schema=sf_config_mock['schema'], table=tmp_territories_table_name ) assert execute_query and expected_table_name in execute_query def test_fill_territories_temp_table(mocker, sf_config_mock): """Test fill_territories_temp_table in snowflake executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) executor = snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock) tmp_territories_table_name = config.snowflake_table_names[ 'territories_temp'] executor.fill_territories_temp_table( tmp_territories_table_name, ows_territories.TEST_TERRITORIES) execute_query, query_params = _find_execute_query_in_mock_calls( mock_connect.mock_calls) expected_table_name = '{db}.{schema}.{table}'.format( db=sf_config_mock['db'], schema=sf_config_mock['schema'], table=tmp_territories_table_name ) assert execute_query assert expected_table_name in execute_query assert query_params == ows_territories.TEST_TERRITORIES_DB_PARAMS def test_create_fact_conflict_temp_table(mocker, sf_config_mock): """Test create_fact_conflict_temp_table function in snowflake_executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) fact_conflict_temp_table = config.snowflake_table_names.get( 'fact_conflict_temp').format(datestamp='20170721') snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).create_fact_conflict_temp_table( fact_conflict_temp_table) execute_query, _ = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert fact_conflict_temp_table in execute_query else: assert False def test_insert_overwrite_into_fact_conflict_temp_table( mocker, sf_config_mock): """Test insert_overwrite_into_fact_conflict_temp_table.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) executor = snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock) art_relations_db = config.art_relations_db art_relations_schema = config.art_relations_schema asset_report_schema = config.asset_report_schema registry_schema = config.registry_schema fact_conflict_temp_table_name = config.snowflake_table_names[ 'fact_conflict_temp'].format(datestamp='20170721') youtube_asset_conflict_by_territory_temp_table_name = \ config.snowflake_table_names[ 'youtube_asset_conflict_by_territory_temp'].format( datestamp='20170721') staging_raw_youtube_asset_report_table_name = config.snowflake_table_names[ 'staging_raw_youtube_asset_report'] registry_table_name = config.snowflake_table_names['registry'] track_table_name = config.snowflake_table_names['track'] releases_table_name = config.snowflake_table_names['releases'] artist_info_table_name = config.snowflake_table_names['artist_info'] territory_standard = config.territory_standard asset_type = config.sound_recording_asset_type time_zone = config.time_zone orchard_account = config.accounts['ORCHARD']['file_label'] executor.insert_overwrite_into_fact_conflict_temp_table( art_relations_db, art_relations_schema, asset_report_schema, registry_schema, fact_conflict_temp_table_name, youtube_asset_conflict_by_territory_temp_table_name, staging_raw_youtube_asset_report_table_name, registry_table_name, track_table_name, releases_table_name, artist_info_table_name, territory_standard, asset_type, time_zone, orchard_account) execute_query, query_params = _find_execute_query_in_mock_calls( mock_connect.mock_calls) expected_table_name = db_schema_table.format( db=sf_config_mock['db'], schema=sf_config_mock['schema'], table=fact_conflict_temp_table_name ) assert execute_query assert 'INSERT OVERWRITE INTO {expected_table_name}'.format( expected_table_name=expected_table_name) in execute_query assert 'FROM ' in execute_query assert 'JOIN ' + db_schema_table.format( db=sf_config_mock['db'], schema=asset_report_schema, table=staging_raw_youtube_asset_report_table_name ) in execute_query assert 'LEFT JOIN ' + db_schema_table.format( db=sf_config_mock['db'], schema=registry_schema, table=registry_table_name ) in execute_query assert 'JOIN ' + db_schema_table.format( db=art_relations_db, schema=art_relations_schema, table='track' ) in execute_query assert 'JOIN ' + db_schema_table.format( db=art_relations_db, schema=art_relations_schema, table='releases' ) in execute_query assert 'JOIN ' + db_schema_table.format( db=art_relations_db, schema=art_relations_schema, table='artist_info' ) in execute_query assert 'WHERE ' in execute_query # temporarily turned off, need to be rewritten after fixing query # assert 'COALESCE(yar.asset_type, yact.asset_type)' in execute_query assert 'AND asset_id IS NOT NULL' in execute_query assert ('AND CONTAINS(UPPER(yar.filename), ' 'UPPER(yact.content_owner))') in execute_query assert 'AND isrc_final IS NOT NULL' in execute_query assert 'QUALIFY' in execute_query assert 'ROW_NUMBER() OVER' in execute_query assert 'ISO_3166_1_2016' in query_params.get('territory_standard') assert 'SOUND_RECORDING' in query_params.get('asset_type') assert 'America/New_York' in query_params.get('time_zone') assert 'ORCH' in query_params.get('orchard_account') def test_check_size_of_fact_conflict_temp_table(mocker, sf_config_mock): """Test check_size_of_fact_conflict_temp_table.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) fact_conflict_temp_table = config.snowflake_table_names.get( 'fact_conflict_temp').format(datestamp='20170721') snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).check_size_of_fact_conflict_temp_table( fact_conflict_temp_table) execute_query, _ = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert fact_conflict_temp_table in execute_query else: assert False def test_check_number_of_unresolved_conflicts(mocker, sf_config_mock): """Test check_number_of_unresolved_conflicts.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) fact_conflict_table = config.snowflake_table_names.get('fact_conflict') snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).check_number_of_unresolved_conflicts( fact_conflict_table) execute_query, _ = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert fact_conflict_table in execute_query else: assert False def test_create_youtube_asset_conflict_by_territory_table( mocker, sf_config_mock): """Test create youtube_asset_conflict_by_territory table.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) yact_table = config.snowflake_table_names.get( 'youtube_asset_conflict_by_territory_temp').format( datestamp='20170101') (snowflake_executor. YouTubeAssetConflictSFExecutor(sf_config_mock). create_youtube_asset_conflict_by_territory_table(yact_table)) execute_query, _ = _find_execute_query_in_mock_calls( mock_connect.mock_calls) assert execute_query and yact_table in execute_query def test_insert_into_fact_conflict_table(mocker, sf_config_mock): """Test insert_into_fact_conflict_table in snowflake_executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) fact_conflict_table = config.snowflake_table_names.get('fact_conflict') fact_conflict_temp_table = config.snowflake_table_names.get( 'fact_conflict_temp').format(datestamp='20170101') snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).insert_into_fact_conflict_table( fact_conflict_table, fact_conflict_temp_table) execute_query, _ = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert fact_conflict_table in execute_query assert fact_conflict_temp_table in execute_query else: assert False def test_update_fact_conflict_yt_recent_daily_average(mocker, sf_config_mock): """Test update_fact_conflict_yt_recent_daily_average snowflake_executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) fact_conflict_table = config.snowflake_table_names.get('fact_conflict') fact_conflict_temp_table = config.snowflake_table_names.get( 'fact_conflict_temp').format(datestamp='20170101') snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).update_fact_conflict_yt_recent_daily_average( fact_conflict_table, fact_conflict_temp_table) execute_query, _ = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert fact_conflict_table in execute_query assert fact_conflict_temp_table in execute_query else: assert False def test_update_fact_conflict_views_in_conflict(mocker, sf_config_mock): """Test update_fact_conflict_views_in_conflict snowflake_executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) fact_conflict_table = config.snowflake_table_names.get('fact_conflict') fact_conflict_temp_table = config.snowflake_table_names.get( 'fact_conflict_temp').format(datestamp='20211230') snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).update_fact_conflict_views_in_conflict( fact_conflict_table, fact_conflict_temp_table) execute_query, _ = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert fact_conflict_table in execute_query assert fact_conflict_temp_table in execute_query else: assert False def test_update_fact_conflict_resolved_datetime(mocker, sf_config_mock): """Test update_fact_conflict_resolved_datetime snowflake_executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) fact_conflict_table = config.snowflake_table_names.get('fact_conflict') fact_conflict_temp_table = config.snowflake_table_names.get( 'fact_conflict_temp').format(datestamp='20170101') time_zone = config.time_zone snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).update_fact_conflict_resolved_datetime( fact_conflict_table, fact_conflict_temp_table, time_zone) execute_query, query_params = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert fact_conflict_table in execute_query assert fact_conflict_temp_table in execute_query else: assert False if query_params: assert time_zone in query_params.get('time_zone') else: assert False def test_reset_es_indexed_for_partially_resolved_conflicts( mocker, sf_config_mock): """Test update fact_conflict es_indexed snowflake executor.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) fact_conflict_table = config.snowflake_table_names.get('fact_conflict') time_zone = config.time_zone snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock).reset_es_indexed_for_partially_resolved_conflicts( fact_conflict_table, time_zone) execute_query, query_params = _find_execute_query_in_mock_calls( mock_connect.mock_calls) if execute_query: assert fact_conflict_table in execute_query else: assert False if query_params: assert time_zone in query_params.get('time_zone') else: assert False def test_fill_youtube_asset_conflict_by_territory_table( mocker, sf_config_mock): """Test fill_youtube_asset_conflict_by_territory_table.""" mock_connect = MagicMock() mocker.patch('snowflake.connector.connect', return_value=mock_connect) executor = snowflake_executor.YouTubeAssetConflictSFExecutor( sf_config_mock) tmp_territories_table_name = config.snowflake_table_names[ 'territories_temp'] yact_table = config.snowflake_table_names.get( 'youtube_asset_conflict_by_territory_temp').format( datestamp='20170101') staging_raw_table = config.snowflake_table_names['staging_raw'] executor.fill_youtube_asset_conflict_by_territory_table( yact_table, staging_raw_table, tmp_territories_table_name) execute_query, query_params = _find_execute_query_in_mock_calls( mock_connect.mock_calls) expected_territories_table_name = '{db}.{schema}.{table}'.format( db=sf_config_mock['db'], schema=sf_config_mock['schema'], table=tmp_territories_table_name ) expected_yact_table_name = '{db}.{schema}.{table}'.format( db=sf_config_mock['db'], schema=sf_config_mock['schema'], table=yact_table ) expected_staging_raw_table_name = '{db}.{schema}.{table}'.format( db=sf_config_mock['db'], schema=sf_config_mock['schema'], table=staging_raw_table ) assert execute_query assert expected_territories_table_name in execute_query assert expected_yact_table_name in execute_query assert expected_staging_raw_table_name in execute_query