"""Tests for Proper Changed Releases flow's tasks.""" import os from unittest.mock import MagicMock from unittest.mock import patch from botocore.exceptions import ClientError from freezegun import freeze_time from feed_sender.flows import tasks as base_tasks from feed_sender.flows.proper_changed_releases import config from feed_sender.flows.proper_changed_releases import tasks @freeze_time('2016-01-01 09:50:00') @patch('feed_sender.flows.proper_changed_releases.tasks.bootstrap_util') def test_bootstrap_reload_false(mock_bootstrap_util): """Test bootstrap with reload false.""" mock_activity = MagicMock() mock_activity.logger.info.return_value = None # Testing reload when feed is already sent. mock_bootstrap_util.feed_already_sent.return_value = True resp = tasks.bootstrap(mock_activity, '2016-01-01 09:50:00', False) assert resp == {'stop': True, 'message': ( 'Feed for 2016-01-01 has already been sent.')} # Testing reload when feed is has not yet been sent. mock_bootstrap_util.feed_already_sent.return_value = False resp = tasks.bootstrap(mock_activity, '2016-01-01 09:50:00', False) config.CHANGED_RELEASES_S3_PATH.format( date='2016-01-01') env = os.getenv('Environment') assert resp == { 'context_date': '2016-01-01', 'feed_name': 'proper_changed_releases', 'cutoff': '2016-01-01 09:50:00', 'changed_releases_s3_path': ( 's3://{}-feed-sender/proper/outgoing/2016-01-01/changed/').format( env), 'changed_releases_filename': 'essential_changes_2016_01_01.csv', 'changed_releases_sftp_path': ('{remote_sftp_path}').format( remote_sftp_path=config.SFTP_REMOTE_PATH), 'release_ids_filename': 'changed_release_ids_2016-01-01.txt', 'filenames_for_sftp_activity': ['essential_changes_2016_01_01.csv'], 'job_ids_filename': 'changed_job_ids_2016-01-01.txt' } @freeze_time('2016-01-01 09:50:00') @patch('feed_sender.flows.proper_changed_releases.tasks.bootstrap_util') @patch('feed_sender.flows.proper_changed_releases.tasks.helpers') def test_bootstrap_reload_true(mock_helpers, mock_bootstrap_util): """Test bootstrap with reload true.""" feed_name = 'proper_changed_releases' context_date = '2016-01-01' cutoff = '2016-01-01 09:50:00' mock_activity = MagicMock() mock_helpers.delete_status.return_value = None mock_activity.logger.info.return_value = None # Testing reload when feed is already sent. mock_bootstrap_util.feed_already_sent.return_value = True resp = tasks.bootstrap(mock_activity, None, True) assert resp == {'stop': True, 'message': ( 'Feed for 2016-01-01 has already been sent.')} # Testing reload when feed is has not yet been sent. mock_bootstrap_util.feed_already_sent.return_value = False resp = tasks.bootstrap(mock_activity, cutoff, True) expected = config.CHANGED_RELEASES_S3_PATH.format( date='2016-01-01') assert resp == { 'context_date': context_date, 'feed_name': 'proper_changed_releases', 'cutoff': cutoff, 'changed_releases_s3_path': expected, 'changed_releases_filename': 'essential_changes_2016_01_01.csv', 'changed_releases_sftp_path': ('{remote_sftp_path}').format( remote_sftp_path=config.SFTP_REMOTE_PATH), 'release_ids_filename': 'changed_release_ids_2016-01-01.txt', 'filenames_for_sftp_activity': ['essential_changes_2016_01_01.csv'], 'job_ids_filename': 'changed_job_ids_2016-01-01.txt' } mock_helpers.delete_status.assert_any_call(feed_name, context_date) @freeze_time('2016-01-01 09:50:00') @patch('feed_sender.flows.proper_changed_releases.tasks.bootstrap_util') def test_bootstrap_reload_false_no_context_date(mock_bootstrap_util): """Test bootstrap with reload true and no context_date.""" mock_activity = MagicMock() mock_activity.logger.info.return_value = None mock_bootstrap_util.feed_already_sent.return_value = True resp = tasks.bootstrap(mock_activity, None, False) assert resp == {'stop': True, 'message': ( 'Feed for 2016-01-01 has already been sent.')} mock_bootstrap_util.feed_already_sent.return_value = False resp = tasks.bootstrap(mock_activity, cutoff=None, reload=False) expected = config.CHANGED_RELEASES_S3_PATH.format( date='2016-01-01') assert resp.get('context_date') == '2016-01-01' assert resp.get('changed_releases_s3_path') == expected @freeze_time('2016-01-01 09:50:00') @patch('feed_sender.flows.proper_changed_releases.tasks.bootstrap_util') def test_bootstrap_no_cutoff(mock_bootstrap_util): """Test bootstrap with no cutoff. (For a feed not sent and reload is False). """ cutoff = '2016-01-01 09:50:00' mock_activity = MagicMock() mock_activity.logger.info.return_value = None # Feed has not yet been sent. mock_bootstrap_util.feed_already_sent.return_value = False resp = tasks.bootstrap(mock_activity, reload=False) expected = config.CHANGED_RELEASES_S3_PATH.format( date='2016-01-01') assert resp == { 'context_date': '2016-01-01', 'feed_name': 'proper_changed_releases', 'cutoff': cutoff, 'changed_releases_s3_path': expected, 'changed_releases_filename': 'essential_changes_2016_01_01.csv', 'changed_releases_sftp_path': ('{remote_sftp_path}').format( remote_sftp_path=config.SFTP_REMOTE_PATH), 'release_ids_filename': 'changed_release_ids_2016-01-01.txt', 'filenames_for_sftp_activity': ['essential_changes_2016_01_01.csv'], 'job_ids_filename': 'changed_job_ids_2016-01-01.txt' } @patch('feed_sender.flows.tasks.helpers') def test_set_status(mock_helpers): """Test set_status.""" mock_helpers.STATUS_SENT = 'SENT' mock_activity = MagicMock() mock_activity.logger.info.return_value = None base_tasks.set_status( mock_activity, 'proper_changed_releases', '2016-01-01', 'SENT') mock_helpers.delete_status.return_value = None mock_helpers.set_status.assert_any_call( 'proper_changed_releases', '2016-01-01', 'SENT') mock_activity.logger.info.assert_any_call( 'Workflow proper_changed_releases on 2016-01-01 has completed.') @patch('feed_sender.flows.proper_changed_releases.tasks.Proper') def test_generate_feed_for_changed_releases_success(mock_proper): """Test success for generate_feed_for_changed_releases task.""" mock_activity = MagicMock() mock_proper.save_release_ids_to_s3.return_value = True squashed_records = MagicMock(records=[1, 2, 3]) mock_proper.return_value.squash_records.return_value = squashed_records context = tasks.generate_feed_for_changed_releases( mock_activity, '2000-06-30 09:50:00', 's3://test', 'abc.csv', 'release_ids.txt', False, 'job_ids.txt') assert context == {'success': True} mock_proper.assert_called_with('2000-06-30 09:50:00') assert mock_proper.return_value.connect_to_sql.called mock_proper.return_value.convert_releases_to_csv.assert_called_with( squashed_records, 's3://test', 'abc.csv') assert mock_proper.return_value.close_sql.called @patch('feed_sender.util.proper_common.smart_open') @patch('feed_sender.flows.proper_changed_releases.tasks.Proper') def test_generate_feed_for_changed_releases_success_from_queue( mock_proper, mock_open, mock_smart_open_obj): """Test success for generate_feed_for_changed_releases task.""" mock_activity = MagicMock() mock_proper.save_release_ids_to_s3.return_value = True squashed_records = MagicMock(records=[1, 2, 3]) mock_proper.return_value.squash_records.return_value = squashed_records mock_open.smart_open.return_value = mock_smart_open_obj mock_proper.return_value.fetch_products.return_value = {} context = tasks.generate_feed_for_changed_releases( mock_activity, '2000-06-30 09:50:00', 's3://test', 'abc.csv', 'release_ids.txt', True, 'job_ids.txt') assert context == {'success': True} mock_proper.assert_called_with('2000-06-30 09:50:00') assert mock_proper.return_value.connect_to_sql.called assert mock_proper.fetch_products mock_proper.return_value.convert_releases_to_csv.assert_called_with( squashed_records, 's3://test', 'abc.csv') assert mock_proper.return_value.close_sql.called @patch('feed_sender.flows.proper_changed_releases.tasks.proper_common') @patch('feed_sender.flows.proper_changed_releases.tasks.job_status.update') def test_update_status_from_job_ids_file( mock_job_status_sql_call, mock_proper_common): """Test if _job_ids are updated in database.""" mock_job_status_sql_call.return_value = True mock_proper_common.get_ids_from_s3.return_value = { 'release_ids': [11111, 22222, 33333], 'job_ids': [1, 2, 3]} update_result = tasks.update_job_status( MagicMock(), 's3://abc', 'file.csv', 'encoding') assert mock_job_status_sql_call.call_count == 3 assert update_result == {'success': True} @patch('feed_sender.flows.proper_changed_releases.tasks.proper_common') @patch('feed_sender.flows.proper_changed_releases.tasks.job_status.update') def test_update_status_failed_from_job_ids_file( mock_job_status_update, mock_proper_common): """Test if _job_ids are updated in database.""" mock_proper_common.get_ids_from_s3.return_value = { 'release_ids': [], 'job_ids': []} update_result = tasks.update_job_status( MagicMock(), 's3://abc', 'file.csv', 'encoding') assert not mock_job_status_update.called assert update_result == {'success': True} @patch('feed_sender.flows.proper_changed_releases.tasks.mysql') @patch('feed_sender.flows.proper_changed_releases.tasks.proper_common') def test_obtain_jobs(mock_proper_common, mock_mysql): """Test when there's messages in the SQS Queue.""" mock_mysql.return_value = MagicMock() mock_proper_common.get_products_from_queue.return_value = { 'release_ids': [1, 2, 3], 'job_ids': [11111, 22222, 33333]} tasks.obtain_jobs(MagicMock(), 'mock_s3_path', 'mock_job_ids_filename') mock_proper_common.save_job_ids_to_s3.assert_called_with( 'mock_s3_path', 'mock_job_ids_filename', [11111, 22222, 33333]) @patch('feed_sender.flows.proper_changed_releases.tasks.mysql') @patch('feed_sender.flows.proper_changed_releases.tasks.proper_common') def test_obtain_jobs_no_jobs(mock_proper_common, mock_mysql): """Test when there's messages in the SQS Queue.""" mock_mysql.get_dd_db_connection_pymysql.return_value = MagicMock() mock_proper_common.get_products_from_queue.return_value = None result = tasks.obtain_jobs(MagicMock(), MagicMock(), MagicMock()) assert not mock_proper_common.save_job_ids_to_s3.called assert result == {'error': 'No jobs available from queue.'} @patch('feed_sender.flows.proper_changed_releases.tasks.mysql') @patch('feed_sender.flows.proper_changed_releases.tasks.proper_common') def test_obtain_jobs_error_queue(mock_proper_common, mock_mysql): """Test when there's messages in the SQS Queue.""" mock_mysql.get_dd_db_connection_pymysql.return_value = MagicMock() mock_proper_common.get_products_from_queue.side_effect = ClientError( {'Error': {'Code': 1, 'Message': 'Mock Client Error'}}, 'MockOp') result = tasks.obtain_jobs(MagicMock(), MagicMock(), MagicMock()) assert not mock_proper_common.save_job_ids_to_s3.called assert result == { 'error': 'An error occurred (1) when calling the MockOp operation: ' 'Mock Client Error'}