import gzip import json import unittest from datetime import datetime import pytest from db_schema.schemas.slz import ( ContentStatus, UnitOfWork, ActivityStatusEnum, CompletenessStatusEnum, UnitOfWorkPriorityEnum, ContentStatusEnum, ContentMetadataStatusEnum ) from unittest.mock import call from requests import Response from slz_appreciationengine_scrapper.const import RATE_LIMIT_HOUR_HEADER from slz_appreciationengine_scrapper.content_status_service import ContentStatusService from slz_appreciationengine_scrapper.dsp.appreciationengine.clients import MembersAllClient from slz_appreciationengine_scrapper.dsp.entities import MetaDataKey from slz_appreciationengine_scrapper.dsp.exceptions import StreamError from slz_appreciationengine_scrapper.entities import Job def bad_response(): resp_mock = unittest.mock.MagicMock(spec=Response) resp_mock.ok = False resp_mock.url = '' resp_mock.status_code = 400 resp_mock.content = b'Rate limit exceeded' resp_mock.reason = b'Rate limit exceeded' resp_mock.headers = {RATE_LIMIT_HOUR_HEADER: 0} return resp_mock def resp_mock_membersall_us_columbia_partial(): data = [ { 'totalSize': 1000, 'items': [ { 'ID': '1000' }, { 'ID': '1002' }, { 'ID': '1003' }, ] }, { 'totalSize': 1000, 'items': [ { 'ID': '1004' }, { 'ID': '1005' }, { 'ID': '1006' }, ] }, { 'totalSize': 1, 'items': [{ 'ID': '1010' }] } ] for dt in data: resp_mock = unittest.mock.MagicMock(spec=Response) resp_mock.ok = True resp_mock.url = '' resp_mock.content = bytes(json.dumps(dt), encoding='utf8') resp_mock.headers = {RATE_LIMIT_HOUR_HEADER: 0} resp_mock.status_code = 200 yield resp_mock @pytest.mark.integration @unittest.mock.patch( 'slz_appreciationengine_scrapper.dsp.appreciationengine.ae_io.time.sleep', return_value=None ) @unittest.mock.patch('slz_appreciationengine_scrapper.dsp.appreciationengine.clients.get_secret') @unittest.mock.patch('requests.get') @unittest.mock.patch( 'slz_appreciationengine_scrapper.dsp.appreciationengine.clients.datetime', wraps=datetime, now=lambda *args, **kwargs: datetime(year=2021, month=6, day=5, hour=1, minute=1, second=1) ) def test_downloading_membersall_partial( datetime_mock, requests_get, get_secret, patched_time_sleep, params, db, s3_client, ae_mapping ): for bucket in ['bucket-archive-quarantine', 'bucket-decompressed-quarantine']: s3_client.create_bucket(Bucket=bucket) uow_dict = { 'uow_id': 'appreciationengine-20201104-sme-membersall-v1', 'dsp': 'appreciationengine', 'report_type': 'membersall', 'version': 'v1', 'report_date': '2020-11-04', 'licensor': 'sme', 'extension': 'json', 'context': 'US_Columbia', } # Client initialization logger = unittest.mock.Mock() timer = unittest.mock.Mock() client = MembersAllClient(logger=logger) # Mock secrets get_secret.return_value = json.dumps({ "Sony Music US - Columbia": "afa0d1d", }) client.configure(params) job = Job.from_dict(uow_dict) expected_file_content = ''.join( [ '{"ID": "1000"}\n', '{"ID": "1002"}\n', '{"ID": "1003"}\n', '{"ID": "1004"}\n', '{"ID": "1005"}\n', '{"ID": "1006"}\n', '{"ID": "1010"}\n', ] ) # Mock response from source API requests_get.return_value.__enter__.side_effect = resp_mock_membersall_us_columbia_partial() content_status_service = ContentStatusService(logger=logger, db_conn=db) meta, meta_data, io_exception = client.download( job, chunk_size=1, timer=timer, content_status_service=content_status_service ) # read data from mocked S3 buckets path = 'appreciationengine/membersall/v1/report_date=2020-11-04/report_licensor=sme' actual_decompressed = s3_client.get_object( Bucket='bucket-decompressed-quarantine', Key=f'{path}/US_Columbia_20201104_20210605010101.json', )['Body'].read() actual_compressed = s3_client.get_object( Bucket='bucket-archive-quarantine', Key=f'{path}/US_Columbia_20201104_20210605010101.json.gz', )['Body'].read() assert expected_file_content == actual_decompressed.decode('utf8') assert expected_file_content == gzip.decompress(actual_compressed).decode('utf8') assert meta_data[MetaDataKey.LAST_MEMBER_ID.value] == '1010' @pytest.mark.integration @unittest.mock.patch( 'slz_appreciationengine_scrapper.dsp.appreciationengine.ae_io.time.sleep', return_value=None ) @unittest.mock.patch('slz_appreciationengine_scrapper.dsp.appreciationengine.clients.get_secret') @unittest.mock.patch('requests.get') @unittest.mock.patch( 'slz_appreciationengine_scrapper.dsp.appreciationengine.clients.datetime', wraps=datetime, now=lambda *args, **kwargs: datetime(year=2021, month=6, day=5, hour=1, minute=1, second=1) ) def test_downloading_membersall_partial_end_member( datetime_mock, requests_get, get_secret, patched_time_sleep, db, params, s3_client, ae_mapping ): for bucket in ['bucket-archive-quarantine', 'bucket-decompressed-quarantine']: s3_client.create_bucket(Bucket=bucket) uow_dict = { 'uow_id': 'appreciationengine-20201104-sme-membersall-v1', 'dsp': 'appreciationengine', 'report_type': 'membersall', 'version': 'v1', 'report_date': '2020-11-04', 'licensor': 'sme', 'extension': 'json', 'context': 'US_Columbia', } # Client initialization logger = unittest.mock.Mock() timer = unittest.mock.Mock() client = MembersAllClient(logger=logger) # Mock secrets get_secret.return_value = json.dumps({ "Sony Music US - Columbia": "afa0d1d", }) client.configure(params) job = Job.from_dict(uow_dict) expected_file_content = ''.join([ '{"ID": "1000"}\n', '{"ID": "1002"}\n', '{"ID": "1003"}\n', ]) # Mock response from source API requests_get.return_value.__enter__.side_effect = resp_mock_membersall_us_columbia_partial() meta_data = { MetaDataKey.LAST_MEMBER_ID.value: 0, MetaDataKey.END_MEMBER_ID.value: 1002, } content_status_service = ContentStatusService(logger=logger, db_conn=db) meta, meta_data, io_exception = client.download( job, chunk_size=1, timer=timer, meta_data=meta_data, content_status_service=content_status_service ) # read data from mocked S3 buckets path = 'appreciationengine/membersall/v1/report_date=2020-11-04/report_licensor=sme' actual_decompressed = s3_client.get_object( Bucket='bucket-decompressed-quarantine', Key=f'{path}/US_Columbia_20201104_20210605010101.json', )['Body'].read() actual_compressed = s3_client.get_object( Bucket='bucket-archive-quarantine', Key=f'{path}/US_Columbia_20201104_20210605010101.json.gz', )['Body'].read() assert expected_file_content == actual_decompressed.decode('utf8') assert expected_file_content == gzip.decompress(actual_compressed).decode('utf8') assert meta_data[MetaDataKey.LAST_MEMBER_ID.value] == 1002 assert meta_data[MetaDataKey.END_MEMBER_ID.value] == 1002 @pytest.mark.integration @unittest.mock.patch( 'slz_appreciationengine_scrapper.dsp.appreciationengine.ae_io.time.sleep', return_value=None ) @unittest.mock.patch('slz_appreciationengine_scrapper.dsp.appreciationengine.clients.get_secret') @unittest.mock.patch('requests.get') @unittest.mock.patch( 'slz_appreciationengine_scrapper.dsp.appreciationengine.clients.datetime', wraps=datetime, now=lambda *args, **kwargs: datetime(year=2021, month=6, day=5, hour=1, minute=1, second=1) ) def test_downloading_membersall_partial_error_raised( datetime_mock, requests_get, get_secret, patched_time_sleep, db, params, s3_client, ae_mapping ): for bucket in ['bucket-archive-quarantine', 'bucket-decompressed-quarantine']: s3_client.create_bucket(Bucket=bucket) uow_dict = { 'uow_id': 'appreciationengine-20201104-sme-membersall-v1', 'dsp': 'appreciationengine', 'report_type': 'membersall', 'version': 'v1', 'report_date': '2020-11-04', 'licensor': 'sme', 'extension': 'json', 'context': 'US_Columbia', } # Client initialization logger = unittest.mock.Mock() timer = unittest.mock.Mock() client = MembersAllClient(logger=logger) # Mock secrets get_secret.return_value = json.dumps({ "Sony Music US - Columbia": "afa0d1d", }) client.configure(params) job = Job.from_dict(uow_dict) expected_file_content = ''.join([ '{"ID": "1000"}\n', '{"ID": "1002"}\n', '{"ID": "1003"}\n', ]) return_value = [*resp_mock_membersall_us_columbia_partial()] # after first request we started to receive an error, most common like limits exceeded return_value[1:1] = [bad_response()] * 4 # Mock response from source API requests_get.return_value.__enter__.side_effect = return_value content_status_service = ContentStatusService(logger=logger, db_conn=db) meta, meta_data, io_exception = client.download( job, chunk_size=1, timer=timer, content_status_service=content_status_service ) # read data from mocked S3 buckets path = 'appreciationengine/membersall/v1/report_date=2020-11-04/report_licensor=sme' actual_decompressed = s3_client.get_object( Bucket='bucket-decompressed-quarantine', Key=f'{path}/US_Columbia_20201104_20210605010101.json', )['Body'].read() actual_compressed = s3_client.get_object( Bucket='bucket-archive-quarantine', Key=f'{path}/US_Columbia_20201104_20210605010101.json.gz', )['Body'].read() assert expected_file_content == actual_decompressed.decode('utf8') assert expected_file_content == gzip.decompress(actual_compressed).decode('utf8') assert meta_data[MetaDataKey.LAST_MEMBER_ID.value] == '1003' @pytest.mark.integration @unittest.mock.patch( 'slz_appreciationengine_scrapper.dsp.appreciationengine.ae_io.time.sleep', return_value=None ) @unittest.mock.patch('slz_appreciationengine_scrapper.dsp.appreciationengine.clients.get_secret') @unittest.mock.patch('requests.get') @unittest.mock.patch( 'slz_appreciationengine_scrapper.dsp.appreciationengine.clients.datetime', wraps=datetime, now=lambda *args, **kwargs: datetime(year=2021, month=6, day=5, hour=1, minute=1, second=1) ) def test_downloading_membersall_partial_error_raised_from_the_beginning( datetime_mock, requests_get, get_secret, patched_time_sleep, db, params, s3_client, ae_mapping ): for bucket in ['bucket-archive-quarantine', 'bucket-decompressed-quarantine']: s3_client.create_bucket(Bucket=bucket) uow_dict = { 'uow_id': 'appreciationengine-20201104-sme-membersall-v1', 'dsp': 'appreciationengine', 'report_type': 'membersall', 'version': 'v1', 'report_date': '2020-11-04', 'licensor': 'sme', 'extension': 'json', 'context': 'US_Columbia', } # Client initialization logger = unittest.mock.Mock() timer = unittest.mock.Mock() client = MembersAllClient(logger=logger) # Mock secrets get_secret.return_value = json.dumps({ "Sony Music US - Columbia": "afa0d1d", }) client.configure(params) job = Job.from_dict(uow_dict) # Mock response from source API # error from the very first request, no content downloaded during this iteration requests_get.return_value.__enter__.side_effect = [bad_response()] * 4 meta_data = {MetaDataKey.LAST_MEMBER_ID.value: 0} content_status_service = ContentStatusService(logger=logger, db_conn=db) with pytest.raises(StreamError): client.download( job, chunk_size=1, timer=timer, meta_data=meta_data, content_status_service=content_status_service ) @unittest.mock.patch( 'slz_appreciationengine_scrapper.dsp.appreciationengine.ae_io.time.sleep', return_value=None ) @unittest.mock.patch('slz_appreciationengine_scrapper.dsp.appreciationengine.clients.get_secret') @unittest.mock.patch('requests.get') def test_ae_availability_check_membersall( requests_get, get_secret, patched_time_sleep, params, logger_test, ae_mapping ): uow = { 'uow_id': 'appreciationengine-20201104-sme-membersall-v1', 'dsp': 'appreciationengine', 'report_type': 'membersall', 'version': 'v1', 'report_date': '2020-11-04', 'licensor': 'sme', 'extension': 'json', 'context': 'CenturyMediaRecords', } job = Job.from_dict(uow) cs_service_mock = unittest.mock.Mock() content_status_mock = ContentStatus(meta_data={MetaDataKey.LAST_MEMBER_ID.value: '123'}) cs_service_mock.get_content_licensor_status_record.side_effect = [ content_status_mock, ] resp_mock_base = unittest.mock.MagicMock(spec=Response) resp_mock_base.ok = True dt = {'items': [{'ID': '123'}], 'totalSize': 1} resp_mock_base.content = bytes(json.dumps(dt), encoding='utf8') resp_mock_base.headers = {} resp_mock_base.url = '' resp_mock_base.status_code = 200 resp_mock_base.reason = 'Some text' requests_get.return_value.__enter__.return_value = resp_mock_base # Mock secrets get_secret.return_value = json.dumps({ 'Sony Music - Century Media Records': 'token72739d3', }) client = MembersAllClient(logger_test) client.configure(params) result = client._check_source_is_available(job, cs_service_mock) url = 'https://sme-delphi.theappreciationengine.com/v1.1/members/all' requests_get.assert_has_calls( [ call(url, { 'apiKey': 'token72739d3', 'limit': "0,1", 'after': '123', 'extended': 1, }), ], any_order=True ) assert result == ('CenturyMediaRecords', 0) @pytest.mark.integration @unittest.mock.patch( 'slz_appreciationengine_scrapper.dsp.appreciationengine.ae_io.time.sleep', return_value=None ) @unittest.mock.patch('slz_appreciationengine_scrapper.dsp.appreciationengine.clients.get_secret') @unittest.mock.patch( 'slz_appreciationengine_scrapper.dsp.appreciationengine.clients.datetime', wraps=datetime, now=lambda *args, **kwargs: datetime(year=2021, month=6, day=5, hour=1, minute=1, second=1) ) def test_downloading_membersall_contexts_member_id( datetime_mock, get_secret, patched_time_sleep, db, s3_client, ae_mapping, params, db_licensors, db_reports, now ): uow_1 = UnitOfWork( unit_of_work_code='appreciationengine-20191124-theorchard-membersall-v1', licensor=db_licensors['theorchard'], reprocess_id='1', report=db_reports['membersall'], report_date='2019-11-24', version='v1', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.session.add(uow_1) cs_1 = ContentStatus( context='UK', content_name='uk.txt', content_status=ContentStatusEnum.ACTIVE, created_at=now, failure_count=0, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, meta_data={ "last_member_id": 1, "first_member_id": 100 } ) uow_1.content_statuses.append(cs_1) db.session.add(cs_1) uow_2 = UnitOfWork( unit_of_work_code='appreciationengine-20191124-sme-membersall-v1', licensor=db_licensors['sme'], reprocess_id='2', report=db_reports['membersall'], report_date='2019-11-24', version='v1', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.session.add(uow_2) cs_2 = ContentStatus( context='UK', content_name='uk.txt', content_status=ContentStatusEnum.ACTIVE, created_at=now, failure_count=0, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, meta_data={ "last_member_id": 100, "first_member_id": 1 } ) uow_2.content_statuses.append(cs_2) db.session.add(cs_2) uow_3 = UnitOfWork( unit_of_work_code='appreciationengine-20191125-theorchard-membersall-v1', licensor=db_licensors['theorchard'], reprocess_id='3', report=db_reports['membersall'], report_date='2019-11-25', version='v1', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.session.add(uow_3) cs_3 = ContentStatus( context='UK', content_name='uk.txt', content_status=ContentStatusEnum.ACTIVE, created_at=now, failure_count=0, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, meta_data={ "last_member_id": 1, "first_member_id": 100 } ) uow_3.content_statuses.append(cs_3) db.session.add(cs_3) uow_4 = UnitOfWork( unit_of_work_code='appreciationengine-20191126-sme-membersall-v1', licensor=db_licensors['sme'], reprocess_id='3', report=db_reports['membersall'], report_date='2019-11-26', version='v1', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.session.add(uow_4) cs_4 = ContentStatus( context='UK', content_name='uk.txt', content_status=ContentStatusEnum.ACTIVE, created_at=now, failure_count=0, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, meta_data={ "last_member_id": 100, "first_member_id": 1 } ) uow_4.content_statuses.append(cs_4) db.session.add(cs_4) db.session.commit() logger = unittest.mock.Mock() content_status_service = ContentStatusService(logger=logger, db_conn=db) max_id = content_status_service.get_max_member_id(cs_3) min_id = content_status_service.get_end_member_id(cs_3) assert max_id == 1 assert min_id == None