"""Tests for validate_raw_data_tasks_sf.""" from unittest.mock import MagicMock, call from unittest.mock import patch from garcon_contrib.dynamo_feed_status import \ garcon_feed_status import pytest from feed_ingestion.tasks import validate_raw_data_tasks_sf from feed_ingestion.util.snowflake.errors import JSONParserError from feed_ingestion.util.snowflake.errors import JSONParserLoading _date = {'2017-11-16'} @pytest.fixture def mock_task_status(): """Yield task status.""" task_status_path = ( 'feed_ingestion.tasks.validate_raw_data_tasks_sf.task_status') with patch(task_status_path) as task_status: task_status.is_completed_task.return_value = False task_status.mark_completed_task = MagicMock() yield task_status @pytest.fixture def mock_set_overall_status(): """Yield overall status.""" overall_status_path = ( 'feed_ingestion.tasks.validate_raw_data_tasks_sf.garcon_feed_status.' 'set_overall_status') with patch(overall_status_path) as overall_status: yield overall_status @pytest.fixture def mock_delete_status(): """Yield overall status.""" status_path = ( 'feed_ingestion.tasks.validate_raw_data_tasks_sf.garcon_feed_status.' 'delete_status') with patch(status_path) as delete_status: yield delete_status class FixturesForValidateTasks: """Fixtures for validate_raw_data_tasks_sf.""" @pytest.fixture def context_validate_raw_data(self, aws_config_mock, sf_config_mock): """Return context for validate_raw_data.""" return (dict( activity=MagicMock(), feed_name='spotify', date=_date, temp_staging_raw_table='test_table', key_dir='s3://path', aws=aws_config_mock, error_limit=2, sfdb_params=sf_config_mock, kwargs={'file_pattern': '.*'})) @pytest.fixture def mock_registered_executors(self, mock_sf_executor): """Mock Snowflake Executor.""" re_path = ( 'feed_ingestion.tasks.validate_raw_data_tasks_sf.' 'registered_executors') with patch(re_path) as mock_re: mock_re.get = MagicMock(return_value=mock_sf_executor) yield mock_re @pytest.fixture def mock_registered_executors_fails(self, mock_sf_executor_fails): """Mock Snowflake Executor returns too many errors.""" re_path = ( 'feed_ingestion.tasks.validate_raw_data_tasks_sf.' 'registered_executors') with patch(re_path) as mock_re: mock_re.get = MagicMock(return_value=mock_sf_executor_fails) yield mock_re @pytest.fixture def mock_sentry_util(self): """Mock sentry_util.""" with patch.object( validate_raw_data_tasks_sf, 'sentry_util') as mock: yield mock class TestValidateRawData(FixturesForValidateTasks): """Test validate_raw_data.""" @pytest.fixture def mock_sf_executor(self): """Mock Snowflake Executor.""" sf_executor = MagicMock() sf_executor.return_value.__enter__.return_value.validate_raw_data = ( MagicMock( return_value=[ JSONParserError( error='err', file='test', line=1, character='c', byte_offset=2, category='error', code='err', sql_state='err', column_name='col', row_number=1, row_start_line=1, rejected_record='err') for _ in range(1)])) yield sf_executor @pytest.fixture def mock_sf_executor_fails(self): """Mock Snowflake Executor returns too many errors.""" sf_executor = MagicMock() sf_executor.return_value.__enter__.return_value.validate_raw_data = ( MagicMock( return_value=[ JSONParserError( error='err', file='test', line=1, character='c', byte_offset=2, category='error', code='err', sql_state='err', column_name='col', row_number=1, row_start_line=1, rejected_record='err') for _ in range(3)])) yield sf_executor def test_log_errors( self, context_validate_raw_data, mock_sentry_util, mock_registered_executors, mock_task_status, mock_set_overall_status, monkeypatch): """Should log invalid entries to Sentry.""" monkeypatch.setenv('SENTRY_DSN', 'http://dsn') validate_raw_data_tasks_sf.validate_raw_data( **context_validate_raw_data) assert mock_sentry_util.send_message.call_args_list == [ call('Invalid records found in: test\n' 'Error: err\n' 'Rejected Record: err\n', level='warning'), ] def test_fails( self, context_validate_raw_data, mock_sentry_util, mock_registered_executors_fails, mock_task_status, mock_set_overall_status, monkeypatch): """Should abort flow if there are too many invalid entries.""" monkeypatch.setenv('SENTRY_DSN', 'http://dsn') with pytest.raises(Exception): validate_raw_data_tasks_sf.validate_raw_data( **context_validate_raw_data) call('Invalid records found in: test\n' 'Error: err\n' 'Rejected Record: err\n', force_send_as_warning=True), mock_set_overall_status.assert_called_with( 'spotify', _date, garcon_feed_status.STATUS_NOT_INGESTED) mock_task_status.mark_completed_task.assert_not_called() class TestLoadTempStagingRawTable(FixturesForValidateTasks): """Test load_temp_staging_raw_table.""" @pytest.fixture def mock_sf_executor(self): """Mock Snowflake Executor.""" sf_executor = MagicMock() sf_executor.return_value.__enter__.return_value.\ load_temp_staging_raw_table = ( MagicMock( return_value=[ JSONParserLoading( file='test', status='LOADED', rows_parsed=100, rows_loaded=100, error_limit=1, errors_seen=0, first_error='NULL', first_error_line='NULL', first_error_character='NULL', first_error_column_name='NULL') for _ in range(2)])) yield sf_executor @pytest.fixture def mock_sf_executor_fails(self): """Mock Snowflake Executor returns too many errors.""" sf_executor = MagicMock() sf_executor.return_value.__enter__.return_value.\ load_temp_staging_raw_table = ( MagicMock( return_value=[ JSONParserLoading( file='test', status='LOADED_FAILED', rows_parsed=100, rows_loaded=0, error_limit=1, errors_seen=100, first_error='1', first_error_line='1', first_error_character='c', first_error_column_name='c') for _ in range(2)])) yield sf_executor @pytest.fixture def mock_sf_executor_partially_loaded(self): """Mock Snowflake Executor returns too many errors.""" sf_executor = MagicMock() sf_executor.return_value.__enter__.return_value.\ load_temp_staging_raw_table = ( MagicMock( return_value=[ JSONParserLoading( file='test', status='PARTIALLY_LOADED', rows_parsed=100, rows_loaded=99, error_limit=10, errors_seen=1, first_error='1', first_error_line='1', first_error_character='c', first_error_column_name='c'), JSONParserLoading( file='test', status='LOADED', rows_parsed=100, rows_loaded=100, error_limit=1, errors_seen=0, first_error='NULL', first_error_line='NULL', first_error_character='NULL', first_error_column_name='NULL'), ])) yield sf_executor @pytest.fixture def mock_registered_executors_partially_loaded( self, mock_sf_executor_partially_loaded): """Mock Snowflake Executor returns some errors.""" re_path = ( 'feed_ingestion.tasks.validate_raw_data_tasks_sf.' 'registered_executors') with patch(re_path) as mock_re: mock_re.get = MagicMock( return_value=mock_sf_executor_partially_loaded) yield mock_re def test_load_temp_staging_raw_table_success( self, context_validate_raw_data, mock_sentry_util, mock_registered_executors, mock_task_status, mock_set_overall_status, mock_delete_status, monkeypatch): """Test when all rows loaded successfully.""" monkeypatch.setenv('SENTRY_DSN', 'http://dsn') validate_raw_data_tasks_sf.load_temp_staging_raw_table( **context_validate_raw_data) assert mock_sentry_util.send_message.call_args_list == [ ] mock_set_overall_status.assert_not_called() mock_delete_status.assert_not_called() def test_load_temp_staging_raw_failed_loading( self, context_validate_raw_data, mock_sentry_util, mock_registered_executors_fails, mock_task_status, mock_set_overall_status, mock_delete_status, monkeypatch): """Should abort flow if there are too many invalid entries.""" monkeypatch.setenv('SENTRY_DSN', 'http://dsn') with pytest.raises(Exception): validate_raw_data_tasks_sf.load_temp_staging_raw_table( **context_validate_raw_data) mock_set_overall_status.assert_called_with( 'spotify', _date, garcon_feed_status.STATUS_NOT_INGESTED) mock_delete_status.assert_called_with( 'spotify', _date) assert mock_sentry_util.send_message.call_args_list == [ call('There are some errors in : test\n' 'Loading status: LOADED_FAILED\n' 'Error: 1\n' 'Count of errors: 100\n\n' 'There are some errors in : test\n' 'Loading status: LOADED_FAILED\n' 'Error: 1\n' 'Count of errors: 100' '\n', level='warning'), ] def test_load_temp_staging_raw_table_partially_loaded( self, context_validate_raw_data, mock_sentry_util, mock_registered_executors_partially_loaded, mock_task_status, mock_set_overall_status, mock_delete_status, monkeypatch): """Should log invalid entries to Sentry.""" monkeypatch.setenv('SENTRY_DSN', 'http://dsn') validate_raw_data_tasks_sf.load_temp_staging_raw_table( **context_validate_raw_data) mock_set_overall_status.assert_not_called() mock_delete_status.assert_not_called() assert mock_sentry_util.send_message.call_args_list == [ call('There are some errors in : test\n' 'Loading status: PARTIALLY_LOADED\n' 'Error: 1\n' 'Count of errors: 1\n', level='warning'), ]