# generic worker test cases from unittest import mock import pytest from slz_appreciationengine_scrapper.config import App from slz_appreciationengine_scrapper.dsp.entities import DecompressedFile, ReportMeta, S3Path from slz_appreciationengine_scrapper.dsp.exceptions import StreamNotFoundError from slz_appreciationengine_scrapper.entities import Job from slz_appreciationengine_scrapper.validation_service import DownloadedSizeError from slz_appreciationengine_scrapper.worker.generic import FileWorker as GenericWorker @pytest.fixture def job_mock(): return Job( uow_id='apple-v1', unit_of_work_id=0, dsp='apple', report_type='users', subtype='', version='v1', report_date='2020-12-02', licensor='smejp', extension='tsv', config_bucket='sme/bucket', context='de', context_params=None, job_id='', ) @pytest.fixture def meta_mock(): path_quarantine = S3Path( bucket='quarantine', path='/path', name='report.tsv', ) path = S3Path( bucket='archive', path='/path', name='report.tsv.gz', ) return ReportMeta( actual_size=123, expected_size=123, destination_path=S3Path( bucket='archive', path='/path', name='report.tsv', ), destination_path_quarantine=path_quarantine, destination_corrupted_path=S3Path( bucket='corrupted', path='/path', name='report.tsv', ), destination_decompressed_path_quarantine=None, files=[ DecompressedFile(name='report.tsv', path=path, path_quarantine=path_quarantine), ], ) @pytest.fixture def s3_service_mock(): return mock.Mock() @pytest.fixture def validation_service_mock(): return mock.Mock() @pytest.fixture def content_status_service_mock(): return mock.Mock() @pytest.fixture def dsp_client_mock(): return mock.Mock() @pytest.fixture def sqs_service_mock(): return mock.Mock() @pytest.fixture def params_mock(): return App.empty() @pytest.fixture def worker_mock( s3_service_mock, validation_service_mock, content_status_service_mock, dsp_client_mock, sqs_service_mock, params_mock ): dsp_settings = mock.Mock() dsp_criterias = mock.Mock() return GenericWorker( mock.Mock(), params=params_mock, s3_service=s3_service_mock, validation_service=validation_service_mock, content_status_service=content_status_service_mock, dsp_client=dsp_client_mock, dsp_settings=dsp_settings, dsp_criterias=dsp_criterias, sqs_service=sqs_service_mock, ) def test_success( validation_service_mock, content_status_service_mock, dsp_client_mock, sqs_service_mock, params_mock, worker_mock, job_mock, meta_mock ): dsp_client_mock.get_content_name.return_value = 'content_name' dsp_client_mock.configure.side_effect = None # exception not raised dsp_client_mock.download.return_value = meta_mock content_status_service_mock.start_processing.return_value = True content_status_service_mock.save_active.return_value = True content_status_service_mock.save_status_success = True validation_service_mock.check_size.side_effect = None validation_service_mock.check_against_schema.return_value = (True, {}, None) validation_service_mock.get_file_size_by_path.return_value = 9999 validation_service_mock.check_calc_size.side_effect = None sqs_service_mock.push.side_effect = None params_mock.validator.enable_historical_validation = True params_mock.validator.threshold_min_file_size = 10000 result = worker_mock.process(job_mock) assert result is True dsp_client_mock.get_content_name.assert_called_with(job_mock) sqs_service_mock.push.assert_called() def test_get_content_name_failed( dsp_client_mock, content_status_service_mock, worker_mock, job_mock ): err = Exception("cannot generate content name") dsp_client_mock.get_content_name.side_effect = err content_status_service_mock.save_failed.return_value = True result = worker_mock.process(job_mock) assert result is False dsp_client_mock.get_content_name.assert_called_with(job_mock) content_status_service_mock.save_failed.assert_called_with( job_mock, '', error_type='Exception', reason=str(err), ) def test_save_status_in_progress_failed( dsp_client_mock, meta_mock, content_status_service_mock, worker_mock, job_mock ): content_name = 'content_name' dsp_client_mock.get_content_name.return_value = content_name dsp_client_mock.download.return_value = meta_mock err = Exception("Failed on start") content_status_service_mock.start_processing.side_effect = err result = worker_mock.process(job_mock) assert result is False dsp_client_mock.get_content_name.assert_called_with(job_mock) content_status_service_mock.save_failed.assert_called_with( job_mock, content_name, error_type='Exception', reason=str(err), ) def test_configuration_failed( dsp_client_mock, content_status_service_mock, worker_mock, job_mock, ): content_name = 'content_name' dsp_client_mock.get_content_name.return_value = content_name err = Exception("Failed on configuration") dsp_client_mock.configure.side_effect = err result = worker_mock.process(job_mock) assert result is False dsp_client_mock.get_content_name.assert_called_with(job_mock) content_status_service_mock.save_failed.assert_called_with( job_mock, content_name, error_type='Exception', reason=str(err), ) def test_save_status_active_failed( dsp_client_mock, content_status_service_mock, worker_mock, job_mock, ): content_name = 'content_name' dsp_client_mock.get_content_name.return_value = content_name dsp_client_mock.configure.side_effect = None content_status_service_mock.save_active.return_value = False result = worker_mock.process(job_mock) assert result is False dsp_client_mock.get_content_name.assert_called_with(job_mock) content_status_service_mock.save_missing.assert_called_with(job_mock, content_name) def test_store_to_quarantine_buckets_failed( dsp_client_mock, content_status_service_mock, s3_service_mock, worker_mock, job_mock, ): content_name = 'content_name' dsp_client_mock.get_content_name.return_value = content_name dsp_client_mock.configure.side_effect = None err = Exception("Failed to download") dsp_client_mock.download.side_effect = err archive_path_quarantine = '/quarantine/archive.gz' dsp_client_mock.get_archive_path_quarantine.return_value = archive_path_quarantine archive_path_corrupted = '/corrupted/archive.gz' dsp_client_mock.get_archive_corrupted_path.return_value = archive_path_corrupted result = worker_mock.process(job_mock) assert result is False dsp_client_mock.get_content_name.assert_called_with(job_mock) content_status_service_mock.save_failed.assert_called_with( job_mock, content_name, error_type='Exception', reason=str(err), ) dsp_client_mock.get_archive_path_quarantine.assert_called() dsp_client_mock.get_archive_corrupted_path.assert_called() s3_service_mock.move.assert_called() def test_store_to_quarantine_buckets_failed_stream_not_found( dsp_client_mock, content_status_service_mock, worker_mock, job_mock, ): content_name = 'content_name' dsp_client_mock.get_content_name.return_value = content_name dsp_client_mock.configure.side_effect = None err = StreamNotFoundError("Failed to download") dsp_client_mock.download.side_effect = err result = worker_mock.process(job_mock) assert result is False dsp_client_mock.get_content_name.assert_called_with(job_mock) content_status_service_mock.save_missing.assert_called_with(job_mock, content_name) def test_file_size_validation_failed( dsp_client_mock, meta_mock, content_status_service_mock, validation_service_mock, s3_service_mock, worker_mock, job_mock, ): content_name = 'content_name' dsp_client_mock.get_content_name.return_value = content_name dsp_client_mock.configure.side_effect = None # exception not raised dsp_client_mock.download.return_value = meta_mock content_status_service_mock.start_processing.return_value = True content_status_service_mock.save_active.return_value = True err = DownloadedSizeError("Size error") validation_service_mock.check_size.side_effect = err result = worker_mock.process(job_mock) assert result is False dsp_client_mock.get_content_name.assert_called_with(job_mock) content_status_service_mock.save_failed.assert_called_with( job_mock, content_name, error_type='DownloadedSizeError', reason=str(err), ) s3_service_mock.delete.assert_called() s3_service_mock.delete_many.assert_called() def test_validate_against_schema_validation_failed( dsp_client_mock, meta_mock, content_status_service_mock, validation_service_mock, s3_service_mock, worker_mock, job_mock, ): content_name = 'content_name' dsp_client_mock.get_content_name.return_value = content_name dsp_client_mock.configure.side_effect = None # exception not raised dsp_client_mock.download.return_value = meta_mock content_status_service_mock.start_processing.return_value = True content_status_service_mock.save_active.return_value = True validation_service_mock.check_size.return_value = True validation_result = (False, {'something': 'happened'}, 'traceback') validation_service_mock.check_against_schema.return_value = validation_result result = worker_mock.process(job_mock) assert result is False dsp_client_mock.get_content_name.assert_called_with(job_mock) content_status_service_mock.save_failed.assert_called_with( job_mock, content_name, error_type='ValidationError', reason=mock.ANY, ) s3_service_mock.move.assert_called() s3_service_mock.delete_many.assert_called() def test_validation_failed( dsp_client_mock, meta_mock, content_status_service_mock, validation_service_mock, worker_mock, job_mock, s3_service_mock, ): content_name = 'content_name' dsp_client_mock.get_content_name.return_value = content_name dsp_client_mock.configure.side_effect = None # exception not raised dsp_client_mock.download.return_value = meta_mock content_status_service_mock.start_processing.return_value = True content_status_service_mock.save_active.return_value = True validation_service_mock.check_size.side_effect = None err = Exception("Unexpected error") validation_service_mock.check_against_schema.side_effect = err result = worker_mock.process(job_mock) assert result is False dsp_client_mock.get_content_name.assert_called_with(job_mock) content_status_service_mock.save_failed.assert_called_with( job_mock, content_name, error_type='Exception', reason=mock.ANY, ) s3_service_mock.move.assert_called() s3_service_mock.delete_many.assert_called() def test_move_to_public_buckets_failed( dsp_client_mock, meta_mock, content_status_service_mock, validation_service_mock, sqs_service_mock, s3_service_mock, params_mock, worker_mock, job_mock, ): content_name = 'content_name' dsp_client_mock.get_content_name.return_value = content_name dsp_client_mock.configure.side_effect = None # exception not raised dsp_client_mock.download.return_value = meta_mock content_status_service_mock.start_processing.return_value = True content_status_service_mock.save_active.return_value = True content_status_service_mock.save_status_success = True validation_service_mock.check_size.side_effect = None validation_service_mock.check_against_schema.return_value = (True, {}, None) validation_service_mock.get_file_size_by_path.return_value = 9999 validation_service_mock.check_calc_size.side_effect = None sqs_service_mock.push.side_effect = None err = Exception("Unhandled error") s3_service_mock.move.side_effect = err params_mock.validator.enable_historical_validation = True params_mock.validator.threshold_min_file_size = 10000 result = worker_mock.process(job_mock) assert result is False dsp_client_mock.get_content_name.assert_called_with(job_mock) content_status_service_mock.save_failed.assert_called_with( job_mock, content_name, error_type='Exception', reason=mock.ANY, ) def test_send_notification_failed( dsp_client_mock, meta_mock, content_status_service_mock, validation_service_mock, sqs_service_mock, params_mock, worker_mock, job_mock, ): content_name = 'content_name' dsp_client_mock.get_content_name.return_value = content_name dsp_client_mock.configure.side_effect = None # exception not raised dsp_client_mock.download.return_value = meta_mock content_status_service_mock.start_processing.return_value = True content_status_service_mock.save_active.return_value = True content_status_service_mock.save_status_success = True validation_service_mock.check_size.side_effect = None validation_service_mock.check_against_schema.return_value = (True, {}, None) validation_service_mock.get_file_size_by_path.return_value = 9999 validation_service_mock.check_calc_size.side_effect = None err = Exception("Unhandled error") sqs_service_mock.push.side_effect = err params_mock.validator.enable_historical_validation = True params_mock.validator.threshold_min_file_size = 10000 result = worker_mock.process(job_mock) assert result is False dsp_client_mock.get_content_name.assert_called_with(job_mock) content_status_service_mock.save_failed.assert_called_with( job_mock, content_name, error_type='Exception', reason=mock.ANY, ) @pytest.mark.parametrize( 'dest_path, move_calls, delete_calls', [( S3Path(bucket='archive', path='/path', name='report.tsv'), 2, 0, ), ( None, 1, 1, )] ) def test_func_move_to_public_bucket( dest_path, move_calls, delete_calls, worker_mock, meta_mock, s3_service_mock, ): meta_mock.destination_path = dest_path worker_mock.move_to_public_buckets(meta_mock) assert s3_service_mock.move.call_count == move_calls assert s3_service_mock.delete.call_count == delete_calls