"""Unit tests for tasks of Apple Music Streams Ingestion Workflow.""" from importlib import reload from itertools import product from unittest import mock from unittest.mock import call, MagicMock from unittest.mock import patch import boto3 from botocore.exceptions import ClientError from freezegun import freeze_time from garcon.activity import Activity from garcon_contrib.dynamo_feed_status import garcon_feed_status import pytest from pytest import raises from feed_ingestion import tasks as common_tasks from feed_ingestion.flows.apple_music_streams import config from feed_ingestion.flows.apple_music_streams import tasks from feed_ingestion.flows.apple_music_streams import vendor_accounts from feed_ingestion.flows.apple_music_streams.utils import \ get_fact_analytics_report from tests.testing_utils import raises_optionally _date = '2021-04-15' @pytest.fixture(autouse=True) def mock_check_status(request, monkeypatch): """Mock check_status decorator. Mock so that the decorated function is kept intact and always get called. """ mock = MagicMock() # just return the function without any modifications mock.return_value = lambda f: f monkeypatch.setattr(common_tasks, 'check_status', value=mock) reload(tasks) # redecorate tasks yield monkeypatch.undo() reload(tasks) @pytest.fixture def mock_get_missing_vendors(): """Yield list of missing vendors.""" missing_vendors_path = \ 'feed_ingestion.flows.apple_music_streams.tasks.get_missing_vendors' with patch(missing_vendors_path) as missing_vendors: missing_vendors.return_value = [] yield missing_vendors mock_contexts_config = { 'default': { 'optional': ['80035150'] } } @pytest.fixture def mock_set_missing_files(): """Yield delete status.""" set_missing_files_path = ( 'feed_ingestion.flows.apple_music_streams.tasks.garcon_feed_status.' 'set_missing_files') with patch(set_missing_files_path) as set_missing_files: yield set_missing_files @pytest.fixture def mock_task_status(): """Yield task status.""" task_status_path = ( 'feed_ingestion.flows.apple_music_streams.tasks.task_status') with patch(task_status_path) as task_status: task_status.is_completed_task.return_value = False task_status.mark_completed_task = MagicMock() task_status.delete_newcontexts = MagicMock() task_status.create_contexts = MagicMock() yield task_status @pytest.fixture def mock_delete_status(): """Yield delete status.""" delete_status_path = ( 'feed_ingestion.flows.apple_music_streams.tasks.garcon_feed_status.' 'delete_status') with patch(delete_status_path) as delete_status: yield delete_status @pytest.fixture def mock_get_overall_status(): """Yield overall status.""" overall_status_path = ( 'feed_ingestion.flows.apple_music_streams.tasks.garcon_feed_status.' 'get_overall_status') with patch(overall_status_path) as overall_status: yield overall_status @pytest.fixture def mock_set_overall_status(): """Yield overall status.""" overall_status_path = ( 'feed_ingestion.flows.apple_music_streams.tasks.' 'set_overall_status_enhanced') with patch(overall_status_path) as overall_status: yield overall_status @pytest.fixture def mock_is_completed_report(): """Yield completed status.""" path = ( 'feed_ingestion.util.task_status.' 'is_completed_report') with patch(path) as f: yield f @pytest.fixture def mock_mark_processed_contexts(): """Yield completed status.""" path = ( 'feed_ingestion.util.task_status.' 'mark_report_processed_contexts') with patch(path) as f: yield f @pytest.fixture def mock_get_values(): """Yield get_values.""" path = ( 'feed_ingestion.util.task_status.' 'get_values') with patch(path) as get_values: yield get_values @pytest.fixture def mock_delete_newcontexts(): """Yield delete new contexts.""" path = ( 'feed_ingestion.util.task_status.' 'delete_newcontexts') with patch(path) as f: yield f @pytest.fixture def mock_is_completed_overall_job(): """Yield delete new contexts.""" path = ( 'feed_ingestion.util.task_status.' 'is_completed_overall_job') with patch(path) as f: yield f @pytest.fixture def mock_get_contexts(): """Yield delete new contexts.""" path = ( 'feed_ingestion.util.task_status.' 'get_report_contexts') with patch(path) as f: yield f @pytest.fixture def mock_add_new_context(): """Yield add new contexts.""" path = ( 'feed_ingestion.util.task_status.' 'add_newcontext') with patch(path) as f: yield f @pytest.fixture def mock_update_context_status(): """Yield contexts.""" path = ( 'feed_ingestion.util.task_status.' 'update_report_context_status') with patch(path) as f: yield f @pytest.fixture def mock_sql_loader(): """Yield SQLLoader.""" path = 'feed_ingestion.flows.apple_music_streams.tasks.SQLLoader' with patch(path) as loader_mock: sql_loader_mock = MagicMock() loader_mock.return_value = sql_loader_mock yield sql_loader_mock @pytest.fixture def expected_bootstrap_response(mock_reports_status_names): """Response for bootstrap task.""" return { 'date': '2020-11-11', 'processed_datetime': '2020-11-11T00:00:00', 's3_archive_bucket': ( 's3://dev-cucumbers/AppleMusicStreams/archives/2020-11-11/'), 's3_drop_bucket': ( 's3://dev-feed-drop/feed-drop/AppleMusicStreams/2020-11-11/'), 'reports_status_names': mock_reports_status_names['theorchard'], 'feed_name_for_fact_analytics': 'apple_music_theorchard_amStreams', 'snowflake_error_limit': config.snowflake_error_limit, 'use_s3': None, 'soft_reload': False, 'fact_analytics_kwargs': {'licensor': 'theorchard'}, 'jenkins_config': config.jenkins_config, 'contexts': vendor_accounts.contexts_config } @pytest.fixture def expected_bootstrap_response_after_switch_on_summary( mock_reports_status_names_after_switch_on_summary): """Response for bootstrap task.""" return { 'date': '2021-04-15', 'processed_datetime': '2021-04-15T00:00:00', 's3_archive_bucket': ( 's3://dev-cucumbers/AppleMusicStreams/archives/2021-04-15/'), 's3_drop_bucket': ( 's3://dev-feed-drop/feed-drop/AppleMusicStreams/2021-04-15/'), 'reports_status_names': mock_reports_status_names_after_switch_on_summary['theorchard'], 'feed_name_for_fact_analytics': 'apple_music_theorchard_amStreamsSummary', 'snowflake_error_limit': config.snowflake_error_limit, 'use_s3': None, 'soft_reload': False, 'fact_analytics_kwargs': {'licensor': 'theorchard'}, 'jenkins_config': config.jenkins_config, 'contexts': { 'default': { 'optional': ['80035150'] } } } @freeze_time('2020-11-11') @patch('feed_ingestion.flows.apple_music_streams.tasks.check_report') def test_bootstrap( mock_check_report, expected_bootstrap_response, mock_delete_status, mock_get_overall_status, mock_set_missing_files, mock_delete_newcontexts, mock_task_status): """Check that bootstrap returns expected results.""" context = { 'activity': MagicMock(), 'date': '2020-11-11', 'reload': False, 'reports': None, 'licensor': 'theorchard' } mock_task_status.is_completed_overall_job.return_value = False mock_task_status.get_report_contexts.return_value = {} response = tasks.bootstrap(**context) assert response == expected_bootstrap_response mock_delete_status.assert_not_called() count = 11 assert mock_check_report.call_count == count @freeze_time('2021-04-15') @patch('feed_ingestion.flows.apple_music_streams.tasks.check_report') @patch('feed_ingestion.flows.apple_music_streams.tasks.contexts_config', mock_contexts_config) def test_bootstrap_after_switch_on_summary( mock_check_report, expected_bootstrap_response_after_switch_on_summary, mock_delete_status, mock_get_overall_status, mock_set_missing_files, mock_is_completed_overall_job, mock_delete_newcontexts, mock_get_contexts): """Check that bootstrap returns expected results.""" context = { 'activity': MagicMock(), 'date': '2021-04-15', 'reload': False, 'reports': None, 'licensor': 'theorchard' } response = tasks.bootstrap(**context) assert response == expected_bootstrap_response_after_switch_on_summary mock_delete_status.assert_not_called() count = 11 assert mock_check_report.call_count == count @pytest.mark.parametrize( 'reload, expected_soft_reload, raises', [ pytest.param('soft', True, None), pytest.param('Soft', True, None), pytest.param('', False, None), pytest.param('1', None, ValueError), pytest.param('tru', None, ValueError), ] ) @freeze_time('2021-04-15') @patch.object(tasks, 'check_report') @patch.object(tasks, 'contexts_config', mock_contexts_config) @patch.object(tasks, 'task_status') def test_bootstrap_soft_reload( mock_task_status, mock_check_report, expected_bootstrap_response_after_switch_on_summary, mock_delete_status, mock_get_overall_status, mock_set_missing_files, mock_is_completed_overall_job, mock_delete_newcontexts, mock_get_contexts, reload, expected_soft_reload, raises ): """Check that bootstrap returns expected results.""" context = { 'activity': MagicMock(), 'date': '2021-04-15', 'reload': reload, 'reports': None, 'licensor': 'theorchard' } with raises_optionally(raises): response = tasks.bootstrap(**context) expected_bootstrap_response_after_switch_on_summary[ 'soft_reload'] = expected_soft_reload assert response == expected_bootstrap_response_after_switch_on_summary mock_delete_status.assert_not_called() if expected_soft_reload: calls = [ call( context_date='2021-04-15', feed_name='apple_music_theorchard', ), call( context_date='2021-04-15', feed_name='apple_music_theorchard_amStreamsSummary', ), ] else: calls = [] assert (mock_task_status. soft_reload_update_status_clearing_tasks. call_args_list) == calls if expected_soft_reload: assert (mock_check_report.call_count == 0), \ 'check_report() is not called when soft reload' else: assert mock_check_report.call_count == 11 @freeze_time('2021-04-15') @patch('feed_ingestion.flows.apple_music_streams.tasks.check_report') def test_bootstrap_reload_reports( mock_check_report, expected_bootstrap_response_after_switch_on_summary, mock_delete_status, mock_set_missing_files, mock_is_completed_overall_job, mock_delete_newcontexts, mock_get_contexts): """Check that bootstrap returns expected results.""" context = { 'activity': MagicMock(), 'date': _date, 'reload': 'True', 'reports': 'amStreamsSummary,amContent', 'licensor': 'theorchard' } response = tasks.bootstrap(**context) reports = ['amStreamsSummary', 'amContent'] reports_status_names = { k: v for k, v in expected_bootstrap_response_after_switch_on_summary[ 'reports_status_names'].items() if k in reports} assert response['reports_status_names'] == reports_status_names mock_delete_status.assert_has_calls( [call('_'.join([config.feed_name, 'theorchard', report]), _date) for report in reports], any_order=True) assert mock_check_report.call_count == len(reports) @freeze_time('2021-04-11') @patch('feed_ingestion.flows.apple_music_streams.tasks.check_report') @patch('feed_ingestion.util.task_status.delete_newcontexts') @patch('feed_ingestion.util.task_status.get_report_contexts') def test_bootstrap_reload_deprecated_reports( mock_get_report_contexts, mock_delete_newcontexts, mock_check_report, expected_bootstrap_response, mock_delete_status, mock_set_missing_files): """Check that bootstrap returns expected results.""" context = { 'activity': MagicMock(), 'date': '2021-04-11', 'reload': 'True', 'reports': 'amStreams', 'licensor': 'theorchard' } with raises(ValueError): tasks.bootstrap(**context) @freeze_time(_date) @patch('feed_ingestion.flows.apple_music_streams.tasks.check_report') def test_bootstrap_reload_licensor( mock_check_report, expected_bootstrap_response_after_switch_on_summary, mock_delete_status, mock_set_missing_files, mock_is_completed_overall_job, mock_delete_newcontexts, mock_get_contexts): """Check that bootstrap returns expected results.""" context = { 'activity': MagicMock(), 'date': _date, 'reload': 'True', 'reports': None, 'licensor': 'theorchard' } response = tasks.bootstrap(**context) assert response['reports_status_names'] == \ expected_bootstrap_response_after_switch_on_summary[ 'reports_status_names'] mock_delete_status.assert_has_calls( [call('_'.join([config.feed_name, 'theorchard', report]), _date) for report in expected_bootstrap_response_after_switch_on_summary[ 'reports_status_names']], any_order=True) assert mock_check_report.call_count == 11 @freeze_time('2020-11-11') @patch('feed_ingestion.flows.apple_music_streams.tasks.check_report') def test_bootstrap_if_some_report_is_already_ingested( mock_check_report, expected_bootstrap_response, mock_delete_status, mock_get_overall_status, mock_set_missing_files, mock_is_completed_overall_job, mock_delete_newcontexts, mock_get_contexts): """Check that bootstrap returns expected results.""" def side_effect_status(report, report_feed_name, date): """Return status False for one of the reports.""" if report == 'amStreams': return False return True context = { 'activity': MagicMock(), 'date': '2020-11-11', 'reload': None, 'reports': None, 'licensor': 'theorchard' } mock_check_report.side_effect = side_effect_status response = tasks.bootstrap(**context) reports_status_names = { k: v for k, v in expected_bootstrap_response[ 'reports_status_names'].items() if k != 'amStreams'} assert response['reports_status_names'] == reports_status_names mock_delete_status.assert_not_called() assert mock_check_report.call_count == 11 @freeze_time('2020-11-11') @patch('feed_ingestion.flows.apple_music_streams.tasks.check_report') def test_bootstrap_snowflake_error_limit( mock_check_report, expected_bootstrap_response, mock_delete_status, mock_get_overall_status, mock_set_missing_files, mock_is_completed_overall_job, mock_delete_newcontexts, mock_get_contexts): """Check that bootstrap returns expected results.""" context = { 'activity': MagicMock(), 'date': '2020-11-11', 'reload': None, 'reports': None, 'licensor': 'theorchard', 'snowflake_error_limit': 100 } response = tasks.bootstrap(**context) assert response['reports_status_names'] == expected_bootstrap_response[ 'reports_status_names'] assert response['snowflake_error_limit'] == 100 mock_delete_status.assert_not_called() assert mock_check_report.call_count == 11 @freeze_time('2020-11-11') @patch('feed_ingestion.flows.apple_music_streams.tasks.check_report') def test_bootstrap_snowflake_use_s3( mock_check_report, expected_bootstrap_response, mock_delete_status, mock_get_overall_status, mock_set_missing_files, mock_delete_newcontexts, mock_is_completed_overall_job, mock_get_contexts): """Check that bootstrap returns expected results.""" context = { 'activity': MagicMock(), 'date': '2020-11-11', 'reload': None, 'reports': None, 'licensor': 'theorchard', 'use_s3': 'True' } response = tasks.bootstrap(**context) assert response['reports_status_names'] == expected_bootstrap_response[ 'reports_status_names'] assert response['use_s3'] == 'True' mock_delete_status.assert_not_called() assert mock_check_report.call_count == 11 @freeze_time('2020-11-11') @patch('feed_ingestion.flows.apple_music_streams.tasks.check_report') def test_bootstrap_complated_1( mock_check_report, expected_bootstrap_response, mock_delete_status, mock_get_overall_status, mock_set_missing_files, mock_is_completed_overall_job, mock_delete_newcontexts, mock_get_contexts): context = { 'activity': MagicMock(), 'date': '2020-11-11', 'reload': None, 'reports': None, 'licensor': 'theorchard' } mock_get_overall_status.return_value = garcon_feed_status.STATUS_INGESTED mock_is_completed_overall_job.return_value = True tasks.bootstrap(**context) mock_delete_status.assert_not_called() assert mock_check_report.call_count == 0 @freeze_time('2020-11-11') @patch('feed_ingestion.flows.apple_music_streams.tasks.check_report') def test_bootstrap_completed_2( mock_check_report, expected_bootstrap_response, mock_delete_status, mock_get_overall_status, mock_set_missing_files, mock_is_completed_overall_job, mock_delete_newcontexts, mock_get_contexts): context = { 'activity': MagicMock(), 'date': '2020-11-11', 'reload': None, 'reports': None, 'licensor': 'theorchard' } mock_get_overall_status.return_value = garcon_feed_status.STATUS_INGESTED mock_is_completed_overall_job.return_value = False tasks.bootstrap(**context) mock_delete_status.assert_not_called() assert mock_check_report.call_count == 11 @freeze_time('2020-11-11') @patch('feed_ingestion.flows.apple_music_streams.tasks.check_report') def test_bootstrap_missing_filed_set_to_empty( mock_check_report, expected_bootstrap_response, mock_delete_status, mock_get_overall_status, mock_set_missing_files, mock_delete_newcontexts, mock_is_completed_overall_job, mock_get_contexts): """Check that bootstrap returns expected results.""" context = { 'activity': MagicMock(), 'date': '2020-11-11', 'reload': None, 'reports': 'amStreams,amContent,amShazam,amTotalLibraryAdds', 'licensor': 'theorchard' } tasks.bootstrap(**context) mock_set_missing_files.assert_any_call( 'apple_music_theorchard_amTotalLibraryAdds', '2020-11-11', []) assert mock_set_missing_files.call_count == 4 def test_check_available_reports_if_all_files_are_available( mock_task_status, mock_reports_status_names_after_switch_on_summary, mock_get_missing_vendors, mock_delete_newcontexts): """Test check_available_reports if all reports are downloaded.""" mock_task_status.is_completed_task.return_value = True for licensor in config.licensors: reports_status_names =\ mock_reports_status_names_after_switch_on_summary[licensor] result = tasks.check_available_reports( MagicMock(), _date, reports_status_names, 's3://te/fe', licensor) for report_name, description in reports_status_names.items(): feed_name = description['feed_name'] mock_task_status.is_completed_task.assert_any_call( feed_name, _date, 'reporter_to_s3') assert result['available_reports'] == reports_status_names if licensor in config.active_licensors: assert result['load_fact_analytics'] is True assert result['update_library_reports'] is True def test_check_available_reports_if_all_files_are_not_available( mock_task_status, mock_reports_status_names, mock_get_missing_vendors): """Test check_available_reports if all reports are not downloaded.""" for licensor in config.licensors: reports_status_names = mock_reports_status_names[licensor] result = tasks.check_available_reports( MagicMock(), _date, reports_status_names, 's3://ts/file', licensor) for report_name, description in reports_status_names.items(): feed_name = description['feed_name'] mock_task_status.is_completed_task.assert_any_call( feed_name, _date, 'reporter_to_s3') assert result == {'stop': True} def test_check_available_reports_if_common_files_are_not_available( mock_task_status, mock_reports_status_names, mock_get_missing_vendors): """Test check_available_reports if common reports are not downloaded.""" def side_effect_status(feed_name, date, task_id): """Return False for common reports otherwise True.""" report = feed_name.split('_', 3)[-1] if report in config.common_reports: return False return True mock_task_status.is_completed_task.side_effect = side_effect_status for licensor in config.licensors: reports_status_names = mock_reports_status_names[licensor] result = tasks.check_available_reports( MagicMock(), _date, reports_status_names, 's3://tet/fie', licensor) for report_name, description in reports_status_names.items(): feed_name = description['feed_name'] mock_task_status.is_completed_task.assert_any_call( feed_name, _date, 'reporter_to_s3') assert result == {'stop': True} def test_check_available_reports_if_only_common_files_are_available( mock_task_status, mock_reports_status_names, mock_get_missing_vendors): """Test check_available_reports if only common reports are downloaded.""" def side_effect_status(feed_name, date, task_id): """Return False for common reports otherwise True.""" report = feed_name.split('_', 3)[-1] if report in config.common_reports: return True return False mock_task_status.is_completed_task.side_effect = side_effect_status for licensor in config.licensors: reports_status_names = mock_reports_status_names[licensor] result = tasks.check_available_reports( MagicMock(), _date, reports_status_names, 's3://tes/fie', licensor) for report_name, description in reports_status_names.items(): feed_name = description['feed_name'] mock_task_status.is_completed_task.assert_any_call( feed_name, _date, 'reporter_to_s3') assert result == {'stop': True} def test_check_available_reports_if_fact_analytics_report_is_not_available( mock_task_status, mock_reports_status_names_after_switch_on_summary, mock_get_missing_vendors): """Test check_available_reports if FA report isn't downloaded.""" def side_effect_status(feed_name, date, task_id): """Return False for common reports otherwise True.""" report = feed_name.split('_', 3)[-1] if report == fact_analytics_report: return False return True mock_task_status.is_completed_task.side_effect = side_effect_status for licensor in config.licensors: fact_analytics_report = get_fact_analytics_report(_date, licensor) reports_status_names = \ mock_reports_status_names_after_switch_on_summary[licensor] result = tasks.check_available_reports( MagicMock(), _date, reports_status_names, 's3://tet/le', licensor) for report_name, description in reports_status_names.items(): feed_name = description['feed_name'] mock_task_status.is_completed_task.assert_any_call( feed_name, _date, 'reporter_to_s3') mock_reports_status_names_after_switch_on_summary[ licensor].pop(fact_analytics_report) assert result['available_reports'] == \ mock_reports_status_names_after_switch_on_summary[licensor] assert result['load_fact_analytics'] is False assert result['update_library_reports'] is True def test_check_available_reports_if_reports_to_update_are_not_available( mock_task_status, mock_reports_status_names_after_switch_on_summary, mock_get_missing_vendors): """Test check_available_reports if library reports aren't downloaded.""" def side_effect_status(feed_name, date, task_id): """Return False for reports_to_update reports otherwise True.""" report = feed_name.split('_', 3)[-1] if report in config.reports_to_update: return False return True mock_task_status.is_completed_task.side_effect = side_effect_status for licensor in config.active_licensors: reports_status_names = \ mock_reports_status_names_after_switch_on_summary[licensor] result = tasks.check_available_reports( MagicMock(), _date, reports_status_names, 's3://tt/fe', licensor) for report_name, description in reports_status_names.items(): feed_name = description['feed_name'] mock_task_status.is_completed_task.assert_any_call( feed_name, _date, 'reporter_to_s3') # remove reports_to_update from output for r in config.reports_to_update: mock_reports_status_names_after_switch_on_summary[licensor].pop(r) assert result['available_reports'] == \ mock_reports_status_names_after_switch_on_summary[licensor] assert result['load_fact_analytics'] is True assert result['update_library_reports'] is False def test_check_available_reports_no_vendors( mock_task_status, mock_reports_status_names_after_switch_on_summary, mock_get_missing_vendors): """Test check_available_reports if there are no new vendors.""" def side_effect_status(feed_name, date, task_id): """Return False for reports_to_update reports otherwise True.""" report = feed_name.split('_', 3)[-1] if report in config.reports_to_update: return False return True mock_task_status.get_report_in_progress_contexts.return_value = [] mock_task_status.is_completed_task.side_effect = side_effect_status for licensor in config.active_licensors: reports_status_names = \ mock_reports_status_names_after_switch_on_summary[licensor] result = tasks.check_available_reports( MagicMock(), _date, reports_status_names, 's3://tt/fe', licensor) assert result == tasks.STOP_RESPONSE def run_staging_raw_table( executor_mock, mock_sql_loader, mock_staging_raw_generator, report_name): """Launch load_staging_raw_table for particular report_name.""" date = '2020-11-11' sfdb_params = {'db': 'db', 'schema': 'schema'} processed_datetime = '2020-11-11T00:00:00' execute_query_mock = MagicMock(return_value=None) executor_mock.get_connection.return_value = MagicMock() executor_mock.return_value.execute_query = execute_query_mock executor_mock.return_value.staging_raw_location.return_value = ( 'db', 'schema') for licensor in config.reports[report_name]['licensors']: vendors = mock_staging_raw_generator['vendors'] staging_raw_table = mock_staging_raw_generator['staging_raw_table'] tasks.load_staging_raw_table( MagicMock(), date, processed_datetime, vendors, sfdb_params, staging_raw_table, report_name, 'apple_music_{}'.format(report_name), licensor) execute_query_mock.assert_any_call( mock_sql_loader, 'delete_from_staging_raw', { 'schema': 'schema', 'db': 'db', 'consumer_db': 'db', 'consumer_schema': 'schema', 'date': '2020-11-11', 'staging_raw_table': staging_raw_table, 'vendor_ids': list(vendors.keys()) } ) for report_account, filename in vendors.items(): temp_staging_raw_table = tasks.get_temp_table_name( '2020-11-11', report_name, report_account, licensor) executor_mock.return_value.load_staging_raw_table.assert_any_call( date, staging_raw_table, temp_staging_raw_table=temp_staging_raw_table, report_name=report_name, licensor=licensor, vendor=report_account, processed_datetime='2020-11-11T00:00:00', filename='test_filename', temp_table_amcontent=tasks.get_temp_table_name( '2020-11-11', 'amContent', report_account, licensor), temp_table_amsubreference=tasks.get_temp_table_name( '2020-11-11', 'amSubscriptionReference', report_account, licensor), apple_id_mapping_table=config.apple_id_mapping_tables[ licensor]) @patch( 'feed_ingestion.flows.apple_music_streams.tasks.AppleMusicStreams') def test_staging_raw_streams_table( executor_mock, mock_sql_loader, mock_staging_raw_generator): """Test load_staging_raw_table task with amStreams report.""" run_staging_raw_table( executor_mock, mock_sql_loader, mock_staging_raw_generator, 'amStreams') @patch( 'feed_ingestion.flows.apple_music_streams.tasks.AppleMusicStreams') def test_staging_raw_demographics_table( executor_mock, mock_sql_loader, mock_staging_raw_generator): """Test load_staging_raw_table task with amContentDemographics report.""" run_staging_raw_table( executor_mock, mock_sql_loader, mock_staging_raw_generator, 'amContentDemographics') @patch( 'feed_ingestion.flows.apple_music_streams.tasks.AppleMusicStreams') def test_staging_raw_nonroyalty_streams_table( executor_mock, mock_sql_loader, mock_staging_raw_generator): """Test load_staging_raw_table task with amNonRoyaltyStreams report.""" run_staging_raw_table( executor_mock, mock_sql_loader, mock_staging_raw_generator, 'amNonRoyaltyStreams') @patch( 'feed_ingestion.flows.apple_music_streams.tasks.AppleMusicStreams') def test_staging_raw_libraryevents_table( executor_mock, mock_sql_loader, mock_staging_raw_generator): """Test load_staging_raw_table task with amLibraryEvents report.""" run_staging_raw_table( executor_mock, mock_sql_loader, mock_staging_raw_generator, 'amLibraryEvents') @patch( 'feed_ingestion.flows.apple_music_streams.tasks.AppleMusicStreams') def test_staging_raw_totallibraryadd_table( executor_mock, mock_sql_loader, mock_staging_raw_generator): """Test load_staging_raw_table task with amTotalLibraryAdds report.""" run_staging_raw_table( executor_mock, mock_sql_loader, mock_staging_raw_generator, 'amTotalLibraryAdds') @patch( 'feed_ingestion.flows.apple_music_streams.tasks.AppleMusicStreams') def test_staging_raw_artists_table( executor_mock, mock_sql_loader, mock_staging_raw_generator): """Test load_staging_raw_table task with amArtists report.""" run_staging_raw_table( executor_mock, mock_sql_loader, mock_staging_raw_generator, 'amArtists') @patch( 'feed_ingestion.flows.apple_music_streams.tasks.AppleMusicStreams') def test_staging_raw_playlists_table( executor_mock, mock_sql_loader, mock_staging_raw_generator): """Test load_staging_raw_table task with amPlaylists report.""" run_staging_raw_table( executor_mock, mock_sql_loader, mock_staging_raw_generator, 'amPlaylists') @patch( 'feed_ingestion.flows.apple_music_streams.tasks.AppleMusicStreams') def test_staging_raw_songs_table( executor_mock, mock_sql_loader, mock_staging_raw_generator): """Test load_staging_raw_table task with amSongs report.""" run_staging_raw_table( executor_mock, mock_sql_loader, mock_staging_raw_generator, 'amSongs') @patch('feed_ingestion.flows.apple_music_streams.tasks.sql_loader') @patch( 'feed_ingestion.flows.apple_music_streams.tasks.SnowflakeSQLExecutor') def test_update_staging_raw_library_reports( executor_mock, sql_loader_mock, mock_get_vendors): """Test update_staging_raw_library_reports task.""" date = '2020-11-11' sfdb_params = {'db': 'db', 'schema': 'schema'} execute_query_mock = MagicMock(return_value=None) executor_mock.get_connection.return_value = MagicMock() executor_mock.return_value.execute_query = execute_query_mock for licensor, report in product( config.licensors, config.reports_to_update): tasks.update_staging_raw_library_reports( MagicMock(), date, 'feed_name', report, sfdb_params, licensor) execute_query_mock.assert_any_call( sql_loader_mock, 'update_staging_raw_library_reports', { 'schema': 'schema', 'db': 'db', 'date': '2020-11-11', 'staging_raw_table': config.reports[ report]['staging_raw_table'], 'vendor_ids': vendor_accounts.get_vendors( '2020-11-11', licensor) }) @patch( 'feed_ingestion.flows.apple_music_streams.tasks.AppleMusicStreams') def test_staging_raw_shazam_table( executor_mock, mock_sql_loader, mock_staging_raw_generator): """Test load_staging_raw_table task with amShazam report.""" run_staging_raw_table( executor_mock, mock_sql_loader, mock_staging_raw_generator, 'amShazam') @patch('feed_ingestion.flows.apple_music_streams.tasks.SnowflakeSQLExecutor') def test_update_apple_id_mapping_with_theorchard_licensor(executor_mock): """Test update_apple_id_mapping.""" execute_query_mock = MagicMock(return_value=None) executor_mock.get_connection.return_value = MagicMock() executor_mock.return_value.__enter__. \ return_value.execute_query = execute_query_mock sfdb_params = {'db': 'db', 'schema': 'schema'} query_names = [ '00_ams_insert_new_entries', '01_update_vendor_offer_code', '02_update_vendor_identifier', '03_clean_blank_apple_release_id', '04_clean_blank_orchard_release_id', '05_clean_zero_apple_track_id', '06_clean_zero_orchard_track_id', '07_update_apple_track_id', '08_update_apple_release_id', '09_update_apple_release_id_one_track', '10_update_apple_release_id_from_upc', '10_ams_update_phonofile_apple_release_id', '11_update_apple_track_id_one_track', '12_update_orchard_track_id', '13_update_orchard_release_id', '14_update_season_pass_orchard_release_track_id', '15_update_with_vendor_identifier_mapping', '16_update_ioda_video_mapping_case', '17_update_ioda_tv_mapping_case', '18_ams_update_orchard_release_track_id_from_isrc', '19_update_orchard_release_id_from_manufacturer_upc_in_vid', '20_update_orchard_release_id_from_manufacturer_upc_in_upc' ] tasks.update_apple_id_mapping( Activity(boto3.client('swf', 'us-east-1')), '2020-11-11', sfdb_params, 'theorchard') for query_name in query_names: execute_query_mock.assert_any_call(mock.ANY, query_name, mock.ANY) @patch('feed_ingestion.flows.apple_music_streams.tasks.SnowflakeSQLExecutor') def test_update_apple_id_mapping_with_sme_licensor(executor_mock): """Test update_apple_id_mapping.""" execute_query_mock = MagicMock(return_value=None) executor_mock.get_connection.return_value = MagicMock() executor_mock.return_value.__enter__. \ return_value.execute_query = execute_query_mock sfdb_params = {'db': 'db', 'schema': 'schema'} result = tasks.update_apple_id_mapping( Activity(boto3.client('swf', 'us-east-1')), '2020-11-11', sfdb_params, 'sme') assert result != {'skip_mapping': True} execute_query_mock.assert_any_call( mock.ANY, '23_populate_sony_apple_id_mapping', mock.ANY) @patch('feed_ingestion.flows.apple_music_streams.tasks.SnowflakeSQLExecutor') def test_update_apple_id_mapping_with_skip_mapping(executor_mock): """Test update_apple_id_mapping with skip_mapping = 'True'.""" execute_query_mock = MagicMock(return_value=None) executor_mock.get_connection.return_value = MagicMock() executor_mock.return_value.__enter__. \ return_value.execute_query = execute_query_mock sfdb_params = {'db': 'db', 'schema': 'schema'} result = tasks.update_apple_id_mapping( Activity(boto3.client('swf', 'us-east-1')), '2020-11-11', sfdb_params, 'licensor', 'True') assert result == {'skip_mapping': True} execute_query_mock.assert_not_called() @patch('feed_ingestion.flows.apple_music_streams.tasks.sql_loader') @patch('feed_ingestion.flows.apple_music_streams.tasks.SnowflakeSQLExecutor') def test_drop_temp_table(executor_mock, sql_loader_mock): """Test load_fact_analytics_error task.""" table_name = 'temp_table' sfdb_params = {'db': 'db', 'schema': 'schema'} execute_query_mock = MagicMock(return_value=None) executor_mock.get_connection.return_value = MagicMock() executor_mock.return_value.__enter__. \ return_value.execute_query = execute_query_mock tasks.drop_temp_table( Activity(boto3.client('swf', 'us-east-1')), table_name, sfdb_params) execute_query_mock.assert_any_call( sql_loader_mock, 'drop_temp_table', { 'schema': 'schema', 'db': 'db', 'table_name': table_name } ) def test_check_report_if_report_is_common( mock_get_overall_status, mock_is_completed_report): """Test function check_report return True if report is common.""" mock_is_completed_report.return_value = False response = tasks.check_report( report_name=(set(config.reports) - set(config.common_reports)).pop(), report_feed_name='test_status', date=_date) assert response is True def test_check_report_if_report_is_ingested( mock_get_overall_status, mock_is_completed_report): """Test function check_report return False if report is ingested.""" mock_is_completed_report.return_value = True mock_get_overall_status.return_value = garcon_feed_status.STATUS_INGESTED response = tasks.check_report( report_name=(set(config.reports) - set(config.common_reports)).pop(), report_feed_name='test_status', date=_date) assert response is False def test_set_status_ingested_if_all_reports_are_ingested( mock_set_overall_status, mock_get_overall_status, mock_reports_status_names, mock_get_values, mock_mark_processed_contexts): """Test set_status_to_ingested if reports are ingested.""" def side_effect_status(feed_name, date): """Set all needed overall statuses.""" report = feed_name.split('_', 3)[-1] if report != fact_analytics_report: return garcon_feed_status.STATUS_POPULATED_RAW_TABLE else: return garcon_feed_status.STATUS_INGESTED mock_get_overall_status.side_effect = side_effect_status for licensor in config.licensors: fact_analytics_report = get_fact_analytics_report(_date, licensor) reports_status_names = mock_reports_status_names[licensor] tasks.set_status_ingested( MagicMock(), _date, reports_status_names, licensor) for report_name, description in reports_status_names.items(): feed_name = description['feed_name'] mock_get_overall_status.assert_any_call(feed_name, _date) if report_name != fact_analytics_report: mock_set_overall_status.assert_any_call( feed_name, _date, garcon_feed_status.STATUS_INGESTED, mock.ANY) def test_set_status_ingested_if_not_all_reports_are_ingested( mock_set_overall_status, mock_get_overall_status, mock_reports_status_names, mock_mark_processed_contexts): """Test set_status_to_ingested if not all reports are done.""" def side_effect_status(feed_name, date): """Set overall statuses not avalible.""" report = feed_name.split('_', 3)[-1] if report != fact_analytics_report: return garcon_feed_status.STATUS_NOT_AVAILABLE else: return garcon_feed_status.STATUS_INGESTED mock_get_overall_status.side_effect = side_effect_status for licensor in config.licensors: fact_analytics_report = get_fact_analytics_report(_date, licensor) reports_status_names = mock_reports_status_names[licensor] tasks.set_status_ingested( MagicMock(), _date, reports_status_names, licensor) for report_name, description in reports_status_names.items(): feed_name = description['feed_name'] mock_get_overall_status.assert_any_call(feed_name, _date) mock_set_overall_status.assert_not_called() def test_set_overall_status_ingested_if_all_reports_are_ingested( mock_set_overall_status, mock_get_overall_status, mock_reports_status_names_after_switch_on_summary, mock_task_status): """Test set_overall_status_to_ingested.""" mock_get_overall_status.return_value = garcon_feed_status.STATUS_INGESTED for licensor in config.licensors: tasks.set_overall_status_ingested(MagicMock(), _date, licensor) for report_name, description in \ mock_reports_status_names_after_switch_on_summary[ licensor].items(): if (licensor, report_name) == ('theorchard', 'amSongs'): continue if (licensor == 'altafonte' and report_name in { 'amArtists', 'amPlaylists', 'amSongs', 'amNonRoyaltySummaryStreams'}): continue mock_get_overall_status.assert_any_call( description['feed_name'], _date) mock_set_overall_status.assert_any_call( f'{config.feed_name}_{licensor}', _date, garcon_feed_status.STATUS_INGESTED, mock.ANY) mock_task_status.mark_completed_task.assert_any_call( f'apple_music_{licensor}_amContent', _date, tasks.STAGING_RAW_TASK_ID) def test_set_overall_status_ingested_if_not_all_reports_are_ingested( mock_set_overall_status, mock_get_overall_status, mock_task_status, mock_reports_status_names_after_switch_on_summary): """Test set_overall_status_to_ingested.""" mock_get_overall_status.return_value = \ garcon_feed_status.STATUS_NOT_AVAILABLE for licensor in config.licensors: response = tasks.set_overall_status_ingested( MagicMock(), _date, licensor) assert response == {'stop': True} mock_set_overall_status.assert_not_called() mock_task_status.mark_completed_task.assert_not_called() def test_get_vendors_config_theorchard(monkeypatch): """Test get_vendors_config.""" monkeypatch.setenv('ORCHARD_VENDOR_PROP_FILE', 'orchard_test') monkeypatch.setenv('IODA_VENDOR_PROP_FILE', 'ioda_test') monkeypatch.setenv('SME_VENDOR_PROP_FILE', 'sme_test') result = tasks.get_vendors_config('theorchard', ['80029727', '80028967']) expected_result = { '80029727': { 'VENDOR_ID': '80029727', 'ACCOUNT': 1682}, '80028967': { 'VENDOR_ID': '80028967', 'ACCOUNT': 1526}} assert expected_result == result def test_get_vendors_config_sme(monkeypatch): """Test get_vendors_config.""" monkeypatch.setenv('ORCHARD_VENDOR_PROP_FILE', 'orchard_test') monkeypatch.setenv('IODA_VENDOR_PROP_FILE', 'ioda_test') monkeypatch.setenv('SME_VENDOR_PROP_FILE', 'sme_test') result = tasks.get_vendors_config('sme', ['80026921', '80030469']) expected_result = { '80026921': { 'VENDOR_ID': '80026921', 'ACCOUNT': 4632}, '80030469': { 'VENDOR_ID': '80030469', 'ACCOUNT': 4632}, } assert expected_result == result @patch( 'feed_ingestion.flows.apple_music_streams.tasks.' 'check_nonconcurrent_workflows') def test_check_concurrent_status_if_no_running_executions( mock_check_nonconcurrent_workflows, monkeypatch): """Test check_concurrent_status.""" max_concurrent_executions = {'licensor': 2} monkeypatch.setattr( config, 'max_concurrent_executions', max_concurrent_executions) mock_check_nonconcurrent_workflows.check_concurrent_status.return_value = \ None result = tasks.check_concurrent_status( MagicMock(), 'licensor', 'test') mock_check_nonconcurrent_workflows.check_concurrent_status. \ assert_called_once_with('test', 'licensor', 2) assert result is None @patch('feed_ingestion.flows.apple_music_streams.' 'tasks.garcon_feed_status.set_status') @patch('feed_ingestion.flows.apple_music_streams.tasks.copy_s3_key') def test_grab_drop_files( mock_copy_s3_key, mock_set_feed_status, mock_add_new_context, mock_update_context_status): """Test grab_drop_files with success run.""" result = tasks.grab_drop_files( MagicMock(), 'feed_name', _date, 's3_archive_path/', 's3_download_path/', 'filename', '1') assert {'filename': garcon_feed_status.STATUS_DOWNLOADED} == result mock_set_feed_status.assert_called_once_with( 'feed_name', _date, 'filename', status=garcon_feed_status.STATUS_DOWNLOADED) mock_copy_s3_key.assert_called_once_with( 's3_download_path/filename', 's3_archive_path/filename') @patch('feed_ingestion.flows.apple_music_streams.' 'tasks.garcon_feed_status.set_status') @patch('feed_ingestion.flows.apple_music_streams.tasks.copy_s3_key') def test_grab_drop_files_run_if_there_is_no_file( mock_copy_s3_key, mock_set_feed_status, mock_add_new_context, mock_update_context_status): """Test grab_drop_files if no file.""" mock_copy_s3_key.side_effect = ClientError( {'Error': {'Code': '404', 'Message': 'Not found'}}, 'Not found') result = tasks.grab_drop_files( MagicMock(), 'feed_name', _date, 's3_archive_path/', 's3_download_path/', 'filename', 'some') assert {'filename': garcon_feed_status.STATUS_NOT_AVAILABLE} == result mock_set_feed_status.assert_not_called() mock_copy_s3_key.assert_called_once_with( 's3_download_path/filename', 's3_archive_path/filename') @patch('feed_ingestion.flows.apple_music_streams.' 'tasks.garcon_feed_status.set_status') @patch('feed_ingestion.flows.apple_music_streams.tasks.copy_s3_key') def test_grab_drop_files_run_if_there_is_no_file_for_optional_vendor( mock_copy_s3_key, mock_set_feed_status, mock_add_new_context, mock_update_context_status): """Test grab_drop_files if no file.""" mock_copy_s3_key.side_effect = ClientError( {'Error': {'Code': '404', 'Message': 'Not found'}}, 'Not found') contexts_config = { 'default': { 'optional': ['80035150'] } } result = tasks.grab_drop_files( MagicMock(), 'feed_name', _date, 's3_archive_path/', 's3_download_path/', 'AppleMusic_TotalLibraryAdds_80035150_20231029', 'some', 'TotalLibraryAdds', contexts_config, ) assert result is None mock_set_feed_status.assert_not_called() mock_copy_s3_key.assert_called_once_with( 's3_download_path/AppleMusic_TotalLibraryAdds_80035150_20231029', 's3_archive_path/AppleMusic_TotalLibraryAdds_80035150_20231029') @patch('feed_ingestion.flows.apple_music_streams.tasks.' 'get_list_of_files_and_directories') @patch('feed_ingestion.flows.apple_music_streams.' 'tasks.garcon_feed_status.set_status') @patch('feed_ingestion.flows.apple_music_streams.tasks.copy_s3_key') def test_grab_drop_files_awal( mock_copy_s3_key, mock_set_feed_status, mock_get_list_of_files_and_directories, mock_add_new_context, mock_update_context_status): """Test grab_drop_files_awal with success run.""" file = 'AppleMusic_SummaryStreams.gz' mock_get_list_of_files_and_directories.return_value = [ f'drop_path/{file}', ] result = tasks.grab_drop_files_awal( MagicMock(), 'feed_name', _date, 's3_archive_path/', 'amStreamsSummary', file, 'some') assert {file: garcon_feed_status.STATUS_DOWNLOADED} == result mock_set_feed_status.assert_called_once_with( 'feed_name', _date, 'AppleMusic_SummaryStreams.gz', status=garcon_feed_status.STATUS_DOWNLOADED) mock_copy_s3_key.assert_called_once_with( 's3://prod-orcd-awal-drop/AWAL_TO_ORCH/ANALYTICS_BACKFILL/' 'applemusic_summary_streams/version=1/year=2021/month=202104/' 'day=20210415/AppleMusic_SummaryStreams.gz', 's3_archive_path/AppleMusic_SummaryStreams.gz') @patch('feed_ingestion.flows.apple_music_streams.tasks.' 'get_list_of_files_and_directories') @patch('feed_ingestion.flows.apple_music_streams.' 'tasks.garcon_feed_status.set_status') @patch('feed_ingestion.flows.apple_music_streams.tasks.copy_s3_key') def test_grab_drop_files_awal_run_if_there_is_no_file( mock_copy_s3_key, mock_set_feed_status, mock_get_list_of_files_and_directories, mock_add_new_context, mock_update_context_status): """Test grab_drop_files_awal if no file.""" file = 'AppleMusic_SummaryStreams.gz' mock_get_list_of_files_and_directories.return_value = [] result = tasks.grab_drop_files_awal( MagicMock(), 'feed_name', _date, 's3_archive_path/', 'amStreamsSummary', file, 'some') assert {file: garcon_feed_status.STATUS_NOT_AVAILABLE} == result mock_set_feed_status.assert_not_called() @patch('feed_ingestion.flows.apple_music_streams.' 'tasks.garcon_feed_status.set_status') @patch('feed_ingestion.flows.apple_music_streams.tasks.copy_s3_key') def test_grab_drop_files_add_new_contexts( mock_copy_s3_key, mock_set_feed_status, mock_add_new_context, mock_update_context_status): """Test grab_drop_files and new contexts.""" tasks.grab_drop_files( MagicMock(), 'feed_name', _date, 's3_archive_path/', 's3_download_path/', 'filename', 'some') mock_add_new_context.assert_called_with('feed_name', _date, 'some') def test_jenkins_config(): assert config.jenkins_config['feeds_required_for_jenkins_build'] == [ 'apple_music_theorchard_amContentDemographics', 'apple_music_theorchard_amContainer', 'apple_music_theorchard_amStreamsSummary', 'apple_music_sme_amContentDemographics', 'apple_music_sme_amContainer', 'apple_music_sme_amStreamsSummary'] @patch('feed_ingestion.flows.apple_music_streams.tasks.overall_status_tasks') def test_set_status_ingested_to_fact_analytics_report( mock_set_overall_status, mock_task_status): mock_task_status.get_values.return_value = ['f1'] licensor = 'sme' fact_analytics_report = get_fact_analytics_report(_date, licensor) tasks.set_status_ingested_to_fact_analytics_report( MagicMock(), _date, fact_analytics_report, garcon_feed_status.STATUS_INGESTED) mock_task_status.mark_report_processed_contexts.assert_called_with( fact_analytics_report, _date)