import datetime from datetime import timedelta from unittest import TestCase, mock from dateutil import tz from parameterized import parameterized from slz_config import DSPSettings from sqlalchemy.exc import OperationalError from slz_downloader.config import App from slz_downloader.dsp.entities import DecompressedFile, ReportMeta, S3Path from slz_downloader.entities import Job from slz_downloader.exceptions import GrasDbConnectionError from slz_downloader.validation_service import CalculatedFileSizeError from slz_downloader.worker.apollo import PlaylistTrackPositionWorker from slz_downloader.worker.apollo import Worker as ApolloWorker from slz_downloader.worker.gras import Worker as GrasWorker class TestAmazonWorkerTestCase(TestCase): pass class TestAppleWorkerTestCase(TestCase): pass class TestApolloWorkerTestCase(TestCase): def setUp(self): self.content_name = 'applemusic_playlists.sql' self.job = Job( uow_id='apple-v1', unit_of_work_id=0, dsp='apollo', report_type='applemusic_playlists', subtype='', version='v1', report_date='2020-12-02', licensor='sme', extension='sql', config_bucket='sme/bucket', context='applemusic_playlists', context_params=None, job_id='', ) self.params = App.empty() self.content_status_service = mock.Mock() self.dsp_client = mock.Mock() self.dsp_settings = DSPSettings.empty() self.dsp_settings.validation_thresholds = {'low': 10, 'high': 20} self.dsp_criterias = mock.Mock() self.validation_service = mock.Mock() self.worker = ApolloWorker( logger=mock.Mock(), params=self.params, s3_service=mock.Mock(), validation_service=self.validation_service, content_status_service=self.content_status_service, dsp_client=self.dsp_client, dsp_settings=self.dsp_settings, dsp_criterias=self.dsp_criterias, sqs_service=mock.Mock(), ) def test_todays_download(self): self.dsp_client.get_content_name.return_value = self.content_name self.dsp_client.download.return_value = 100 self.validation_service.validate.return_value = True datetime_mock = mock.Mock(wraps=datetime.datetime) datetime_mock.today.return_value = datetime.datetime(2020, 12, 2) with mock.patch('slz_downloader.worker.apollo.datetime', new=datetime_mock): result = self.worker.process(self.job) self.assertFalse(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.content_status_service.start_processing.assert_called_with(self.job, self.content_name) self.dsp_client.configure.assert_called_with( self.params, validation_service=self.validation_service ) self.content_status_service.save_active.assert_called_with(self.job, self.content_name) self.dsp_client.download.assert_called_with(self.job) self.content_status_service.save_missing.assert_called_with(self.job, self.content_name) def test_configuration_failed(self): self.dsp_client.get_content_name.return_value = self.content_name self.dsp_client.configure.side_effect = Exception("Misconfiguration") result = self.worker.process(self.job) self.assertFalse(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.content_status_service.save_failed.assert_called_with( self.job, self.content_name, error_type='Exception', reason="Misconfiguration", ) def test_outdated_download(self): self.dsp_client.get_content_name.return_value = self.content_name self.dsp_client.download.return_value = 6789 self.validation_service.validate.return_value = True datetime_mock = mock.Mock(wraps=datetime.datetime) datetime_mock.today.return_value = datetime.datetime( 2020, 12, 3, 1, 2, 3, 1234, tzinfo=tz.gettz('America/New_York') ) with mock.patch('slz_downloader.worker.apollo.datetime', new=datetime_mock): result = self.worker.process(self.job) self.assertTrue(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.content_status_service.start_processing.assert_called_with(self.job, self.content_name) self.dsp_client.configure.assert_called_with( self.params, validation_service=self.validation_service ) self.content_status_service.save_active.assert_called_with(self.job, self.content_name) self.dsp_client.download.assert_called_with(self.job) self.content_status_service.save_success.assert_called_with( self.job, self.content_name, self.content_name, 6789, [], ) def test_outdated_download_empty_result(self): self.dsp_client.get_content_name.return_value = self.content_name self.dsp_client.download.return_value = 0 datetime_mock = mock.Mock(wraps=datetime.datetime) datetime_mock.today.return_value = datetime.datetime( 2020, 12, 3, 1, 2, 3, 1234, tzinfo=tz.gettz('America/New_York') ) with mock.patch('slz_downloader.worker.apollo.datetime', new=datetime_mock): result = self.worker.process(self.job) self.assertFalse(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.validation_service.validate.return_value = False self.content_status_service.start_processing.assert_called_with(self.job, self.content_name) self.dsp_client.configure.assert_called_with( self.params, validation_service=self.validation_service ) self.content_status_service.save_active.assert_called_with(self.job, self.content_name) self.dsp_client.download.assert_called_with(self.job) self.validation_service.validate.assert_not_called() self.content_status_service.save_missing.assert_called_with(self.job, self.content_name) def test_outdated_download_failed_validation(self): self.params.apollo.complete_days_offset = 0 self.dsp_client.get_content_name.return_value = self.content_name self.dsp_client.download.return_value = 123 err = CalculatedFileSizeError() self.validation_service.check_filesize_history.side_effect = err datetime_mock = mock.Mock(wraps=datetime.datetime) datetime_mock.today.return_value = datetime.datetime( 2020, 12, 3, 1, 2, 3, 1234, tzinfo=tz.gettz('America/New_York') ) with mock.patch('slz_downloader.worker.apollo.datetime', new=datetime_mock): result = self.worker.process(self.job) self.assertFalse(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.content_status_service.start_processing.assert_called_with(self.job, self.content_name) self.dsp_client.configure.assert_called_with( self.params, validation_service=self.validation_service ) self.content_status_service.save_active.assert_called_with(self.job, self.content_name) self.dsp_client.download.assert_called_with(self.job) self.validation_service.validate.assert_not_called() self.content_status_service.save_on_hold.assert_called_with( self.job, self.content_name, 123 ) def test_outdated_download_missing_offset_5(self): self.params.apollo.complete_days_offset = 5 self.dsp_client.get_content_name.return_value = self.content_name self.dsp_client.download.return_value = 123 err = CalculatedFileSizeError() self.validation_service.check_filesize_history.side_effect = err datetime_mock = mock.Mock(wraps=datetime.datetime) datetime_mock.today.return_value = datetime.datetime( 2020, 12, 3, 1, 2, 3, 1234, tzinfo=tz.gettz('America/New_York') ) with mock.patch('slz_downloader.worker.apollo.datetime', new=datetime_mock): result = self.worker.process(self.job) self.assertFalse(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.content_status_service.start_processing.assert_called_with(self.job, self.content_name) self.dsp_client.configure.assert_called_with( self.params, validation_service=self.validation_service ) self.content_status_service.save_active.assert_called_with(self.job, self.content_name) self.dsp_client.download.assert_called_with(self.job) self.validation_service.validate.assert_not_called() self.content_status_service.save_missing.assert_called_with( self.job, self.content_name, ) def test_outdated_download_failed_validation_offset_5(self): self.params.apollo.complete_days_offset = 5 self.dsp_client.get_content_name.return_value = self.content_name self.dsp_client.download.return_value = 123 err = CalculatedFileSizeError() self.validation_service.check_filesize_history.side_effect = err datetime_mock = mock.Mock(wraps=datetime.datetime) datetime_mock.today.return_value = datetime.datetime( 2020, 12, 8, 1, 2, 3, 1234, tzinfo=tz.gettz('America/New_York') ) with mock.patch('slz_downloader.worker.apollo.datetime', new=datetime_mock): result = self.worker.process(self.job) self.assertFalse(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.content_status_service.start_processing.assert_called_with(self.job, self.content_name) self.dsp_client.configure.assert_called_with( self.params, validation_service=self.validation_service ) self.content_status_service.save_active.assert_called_with(self.job, self.content_name) self.dsp_client.download.assert_called_with(self.job) self.validation_service.validate.assert_not_called() self.content_status_service.save_on_hold.assert_called_with( self.job, self.content_name, 123 ) def test_failed_download(self): class DownloadError(Exception): pass self.dsp_client.get_content_name.return_value = self.content_name self.dsp_client.download.side_effect = DownloadError("Download failed") result = self.worker.process(self.job) self.assertFalse(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.content_status_service.save_failed.assert_called_with( self.job, self.content_name, error_type='DownloadError', reason="Download failed", ) @parameterized.expand( [ ('2020-12-09', None, False), ('2020-12-08', None, False), ('2020-12-07', None, True), ('2020-12-01', timedelta(days=3), True), ('2020-12-06', timedelta(days=1), True), ('2020-12-06', timedelta(days=2), False), ] ) def test_outdated(self, report_date, delta, expect): job = self.job job.report_date = report_date datetime_mock = mock.Mock(wraps=datetime.datetime) datetime_mock.today.return_value = datetime.datetime( 2020, 12, 8, 1, 2, 3, 1234, tzinfo=tz.gettz('America/New_York') ) with mock.patch('slz_downloader.worker.apollo.datetime', new=datetime_mock): result = ApolloWorker.is_outdated_job(job, delta) self.assertEqual(expect, result) class PlaylistTrackPositionWorkerTestCase(TestCase): def setUp(self): self.content_name = 'applemusic_playlist_track_position.sql' self.job = Job( uow_id='apple-v1', unit_of_work_id=0, dsp='apollo', report_type='applemusic_playlist_track_position', subtype='', version='v1', report_date='2020-12-02', licensor='sme', extension='sql', config_bucket='sme/bucket', context='applemusic_playlist_track_position', context_params=None, job_id='', ) self.params = App.empty() self.content_status_service = mock.Mock() self.dsp_client = mock.Mock() self.dsp_settings = DSPSettings.empty() self.dsp_settings.validation_thresholds = {'low': 10, 'high': 20} self.dsp_criterias = mock.Mock() self.validation_service = mock.Mock() self.worker = PlaylistTrackPositionWorker( logger=mock.Mock(), params=self.params, s3_service=mock.Mock(), validation_service=self.validation_service, content_status_service=self.content_status_service, dsp_client=self.dsp_client, dsp_settings=self.dsp_settings, dsp_criterias=self.dsp_criterias, sqs_service=mock.Mock(), ) def test_outdated_playlist_track_position_complete(self): self.content_status_service.get_file_size.return_value = 123 self.content_status_service.get_file_size_history.return_value = [100, 200, 300] self.dsp_client.get_content_name.return_value = self.content_name datetime_mock = mock.Mock(wraps=datetime.datetime) datetime_mock.today.return_value = datetime.datetime( 2020, 12, 3, 1, 2, 3, 1234, tzinfo=tz.gettz('America/New_York') ) with mock.patch('slz_downloader.worker.apollo.datetime', new=datetime_mock): result = self.worker.process(self.job) self.assertTrue(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.content_status_service.start_processing.assert_called_with(self.job, self.content_name) self.dsp_client.configure.assert_called_with( self.params, validation_service=self.validation_service ) self.content_status_service.save_active.assert_called_with(self.job, self.content_name) self.dsp_client.download.assert_not_called() self.validation_service.check_filesize_history.assert_called_with( mock.ANY, 123, [100, 200, 300], 10, 20 ) self.content_status_service.save_success.assert_called_with( self.job, self.content_name, self.content_name, 123, [], ) def test_outdated_playlist_track_position_incomplete_invalid(self): self.content_status_service.get_file_size.return_value = 1000 self.content_status_service.get_file_size_history.return_value = [100, 200, 300] self.dsp_client.get_content_name.return_value = self.content_name err = CalculatedFileSizeError() self.validation_service.check_filesize_history.side_effect = err datetime_mock = mock.Mock(wraps=datetime.datetime) datetime_mock.today.return_value = datetime.datetime( 2020, 12, 3, 1, 2, 3, 1234, tzinfo=tz.gettz('America/New_York') ) self.params.apollo.complete_days_offset = 0 with mock.patch('slz_downloader.worker.apollo.datetime', new=datetime_mock): result = self.worker.process(self.job) self.assertFalse(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.content_status_service.start_processing.assert_called_with(self.job, self.content_name) self.dsp_client.configure.assert_called_with( self.params, validation_service=self.validation_service ) self.content_status_service.save_active.assert_called_with(self.job, self.content_name) self.dsp_client.download.assert_not_called() self.validation_service.check_filesize_history.assert_called_with( mock.ANY, 1000, [100, 200, 300], 10, 20 ) self.content_status_service.save_on_hold.assert_called_with( self.job, self.content_name, 1000 ) def test_outdated_playlist_track_position_incomplete(self): file_size = 0 self.content_status_service.get_file_size.return_value = file_size self.dsp_client.get_content_name.return_value = self.content_name datetime_mock = mock.Mock(wraps=datetime.datetime) datetime_mock.today.return_value = datetime.datetime( 2020, 12, 3, 1, 2, 3, 1234, tzinfo=tz.gettz('America/New_York') ) with mock.patch('slz_downloader.worker.apollo.datetime', new=datetime_mock): result = self.worker.process(self.job) self.assertFalse(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.content_status_service.start_processing.assert_called_with(self.job, self.content_name) self.dsp_client.configure.assert_called_with( self.params, validation_service=self.validation_service ) self.content_status_service.save_active.assert_called_with(self.job, self.content_name) self.dsp_client.download.assert_not_called() self.validation_service.check_filesize_history.assert_not_called() self.content_status_service.save_on_hold.assert_called_with( self.job, self.content_name, file_size ) def test_today_playlist_track_position(self): self.dsp_client.get_content_name.return_value = self.content_name self.dsp_client.download.return_value = 5667 datetime_mock = mock.Mock(wraps=datetime.datetime) datetime_mock.today.return_value = datetime.datetime( 2020, 12, 2, 1, 2, 3, 1234, tzinfo=tz.gettz('America/New_York') ) with mock.patch('slz_downloader.worker.apollo.datetime', new=datetime_mock): result = self.worker.process(self.job) self.assertFalse(result) self.dsp_client.get_content_name.assert_called_with(self.job) self.content_status_service.start_processing.assert_called_with(self.job, self.content_name) self.dsp_client.configure.assert_called_with( self.params, validation_service=self.validation_service ) self.content_status_service.save_active.assert_called_with(self.job, self.content_name) self.dsp_client.download.assert_called_with(self.job) self.content_status_service.update_file_size.assert_called_with(self.job, 5667) self.content_status_service.save_missing.assert_called_with(self.job, self.content_name) class GrasWorkerTestCase(TestCase): def setUp(self): self._job = Job( uow_id='gras-v1', unit_of_work_id=0, dsp='gras', report_type='gras_artist', subtype='', version='v1', report_date='2020-12-02', licensor='sme', extension='csv', config_bucket='sme/bucket', context='de', context_params=None, job_id='', ) path_quarantine = S3Path( bucket='quarantine', path='/path', name='report.tsv', ) path = S3Path( bucket='archive', path='/path', name='report.tsv.gz', ) self._meta = 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), ], ) self._s3_service = mock.Mock() self._validation_service = mock.Mock() self._content_status_service = mock.Mock() self._dsp_client = mock.Mock() self._dsp_settings = mock.Mock() self._dsp_criterias = mock.Mock() self._sqs_service = mock.Mock() self._worker = GrasWorker( mock.Mock(), params=App.empty(), s3_service=self._s3_service, validation_service=self._validation_service, content_status_service=self._content_status_service, dsp_client=self._dsp_client, dsp_settings=self._dsp_settings, dsp_criterias=self._dsp_criterias, sqs_service=self._sqs_service, ) def test_success(self): self._dsp_client.get_content_name.return_value = 'content_name' self._dsp_client.configure.side_effect = None # exception not raised meta = self._meta self._dsp_client.download.return_value = meta self._content_status_service.start_processing.return_value = True self._content_status_service.save_active.return_value = True self._content_status_service.save_status_success = True self._validation_service.check_size.side_effect = None self._validation_service.check_against_schema.return_value = (True, {}, None) self._validation_service.get_file_size_by_path.return_value = 9999 self._validation_service.check_calc_size.side_effect = None self._sqs_service.push.side_effect = None params = App.empty() params.validator.enable_historical_validation = True params.validator.threshold_min_file_size = 10000 self._worker._params = params job = self._job result = self._worker.process(job) self.assertTrue(result) self._dsp_client.get_content_name.assert_called_with(job) def test_store_to_quarantine_buckets_failed(self): content_name = 'content_name' self._dsp_client.get_content_name.return_value = content_name self._dsp_client.configure.side_effect = None err = Exception("Failed to download") self._dsp_client.download.side_effect = err archive_path_quarantine = '/quarantine/archive.gz' self._dsp_client.get_archive_path_quarantine.return_value = archive_path_quarantine archive_path_corrupted = '/corrupted/archive.gz' self._dsp_client.get_archive_corrupted_path.return_value = archive_path_corrupted result = self._worker.process(self._job) self.assertFalse(result) self._dsp_client.get_content_name.assert_called_with(self._job) self._content_status_service.save_failed.assert_called_with( self._job, content_name, error_type='Exception', reason=str(err), ) self._dsp_client.get_archive_path_quarantine.assert_called() self._dsp_client.get_archive_corrupted_path.assert_called() self._s3_service.move.assert_called() def test_store_to_quarantine_buckets_failed_with_no_db_conn(self): content_name = 'content_name' self._dsp_client.get_content_name.return_value = content_name self._dsp_client.configure.side_effect = None err = OperationalError("pyodbc.OperationalError", params={}, orig={}) self._dsp_client.download.side_effect = err self.assertRaises(GrasDbConnectionError, self._worker.process, self._job) self._dsp_client.get_content_name.assert_called_with(self._job) self._content_status_service.save_failed.assert_called_with( self._job, content_name, error_type='OperationalError', reason=str(err), )