from datetime import datetime, timedelta, timezone from unittest import TestCase, mock import pytest from db_schema.common import ActivityStatusEnum, CompletenessStatusEnum from db_schema.schemas import slz from db_schema.schemas.slz import ( ContentMetadataStatusEnum, ContentStatus, ContentStatusEnum, UnitOfWork, UnitOfWorkPriorityEnum, ) from slz_appreciationengine_scrapper.content_status_service import ContentStatusService from slz_appreciationengine_scrapper.dsp.entities import AEDateTimeRange from slz_appreciationengine_scrapper.entities import Job @pytest.fixture def job_mock(uow_mock): return Job( uow_id='appreciationengine-20210825-sme-membersvisittotals-v1', unit_of_work_id=uow_mock.unit_of_work_id, dsp='appreciationengine', report_type='membersvisittotals', subtype=None, version='v1', report_date='2021-08-25', licensor='sme', extension='csv', config_bucket='delphi-configs', context='UK', context_params={}, job_id="test_job_id", ) @pytest.fixture def uow_mock(db, db_licensors, db_reports): now = datetime.now(timezone.utc) - timedelta(days=1) uow = UnitOfWork( unit_of_work_code='appreciationengine-20210825-sme-membersvisittotals-v1', licensor=db_licensors['sme'], report=db_reports['membersvisittotals'], 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=now, last_updated_at=now, created_at=now, ) db.session.add(uow) db.session.commit() return uow @pytest.fixture def cs_mock(db, uow_mock, job_mock): now = datetime.now(timezone.utc) - timedelta(days=1) cs = slz.ContentStatus( unit_of_work_id=uow_mock.unit_of_work_id, content_name='report.txt.gz', context=job_mock.context, content_status=slz.ContentStatusEnum.ACTIVE.value, latest_job_id=job_mock.job_id, content_size=44, created_at=now, last_checked_at=now, failure_count=0, metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) db.session.add(cs) db.session.commit() return cs @pytest.mark.integration @mock.patch('slz_appreciationengine_scrapper.content_status_service.datetime') def test_save_failed(datetime_mock, logger_test, db, job_mock, cs_mock): datetime_mock.now.return_value = datetime(2019, 4, 4, 1) service = ContentStatusService(logger=logger_test, db_conn=db) assert db.session.query(slz.ContentFailureLog).count() == 0 result = service.save_failed( job_mock, 'content_name', 'ValidationError', 'some fields are invalid', ) assert result is True assert cs_mock.content_status == slz.ContentStatusEnum.FAILED assert db.session.query(slz.ContentFailureLog).count() == 1 @pytest.mark.integration @mock.patch('slz_appreciationengine_scrapper.content_status_service.datetime') def test_save_success(datetime_mock, logger_test, db, job_mock, cs_mock): datetime_mock.now.return_value = datetime(2019, 4, 4, 1) service = ContentStatusService(logger=logger_test, db_conn=db) result = service.save_success( job_mock, 'content_name', 'content_name', 444, [], ) assert result is True assert cs_mock.content_status == slz.ContentStatusEnum.COMPLETE assert cs_mock.content_size == 444 @pytest.mark.integration @mock.patch('slz_appreciationengine_scrapper.content_status_service.datetime') def test_save_active(datetime_mock, logger_test, db, job_mock, cs_mock): datetime_mock.now.return_value = datetime(2019, 4, 4, 1) service = ContentStatusService(logger=logger_test, db_conn=db) result = service.save_active( job_mock, 'content_name', ) assert result is True assert cs_mock.content_status == slz.ContentStatusEnum.ACTIVE @pytest.mark.integration @mock.patch('slz_appreciationengine_scrapper.content_status_service.datetime') def test_save_missing(datetime_mock, logger_test, db, job_mock, cs_mock): cs_mock.content_status = slz.ContentStatusEnum.FAILED db.session.commit() assert cs_mock.content_status == slz.ContentStatusEnum.FAILED datetime_mock.now.return_value = datetime(2019, 4, 4, 1) service = ContentStatusService(logger=logger_test, db_conn=db) result = service.save_missing( job_mock, 'content_name', ) assert result is True assert cs_mock.content_status == slz.ContentStatusEnum.MISSING @pytest.mark.integration @mock.patch('slz_appreciationengine_scrapper.content_status_service.datetime') def test_start_processing(datetime_mock, logger_test, db, job_mock, cs_mock): datetime_mock.now.return_value = now = datetime(2019, 4, 4, 1, tzinfo=timezone.utc) service = ContentStatusService(logger=logger_test, db_conn=db) # already completed cs_mock.content_status = slz.ContentStatusEnum.COMPLETE db.session.commit() result = service.start_processing( job_mock, 'content_name', ) assert result is False assert db.session.query(slz.ContentStatus).count() == 1 assert cs_mock.last_checked_at != now # Updated cs_mock.content_status = slz.ContentStatusEnum.FAILED db.session.commit() result = service.start_processing( job_mock, 'content_name', ) assert result is True assert db.session.query(slz.ContentStatus).count() == 1 assert cs_mock.last_checked_at == now # new job_mock.context = 'Alamo' result = service.start_processing( job_mock, 'content_name', ) assert result is True assert db.session.query(slz.ContentStatus).count() == 2 cs = db.session.query(slz.ContentStatus).filter(slz.ContentStatus.context == job_mock.context ).one() assert cs.created_at == now @pytest.mark.integration def test_get_max_member_id_should_return_0_if_no_records(db, db_licensors, db_reports, logger_test): uow_1 = UnitOfWork( unit_of_work_code='appreciationengine-20210825-sme-membersvisittotals-v1', licensor=db_licensors['sme'], report=db_reports['membersvisittotals'], 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-01T01:01:01.377644+00:00', last_updated_at='2019-11-01T01:00:39.985786+00:00', created_at='2019-11-01T01:01:00.599872+00:00', ) cs = ContentStatus( unit_of_work=uow_1, context='US_Columbia', content_name='US_Columbia.txt', content_status=ContentStatusEnum.MISSING, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) db.session.add_all([uow_1, cs]) db.session.commit() service = ContentStatusService( logger=logger_test, db_conn=db, ) max_id = service.get_max_member_id(cs) assert max_id == 0 @pytest.mark.integration def test_get_max_member_id_ok(db, db_licensors, db_reports, logger_test): uow_1 = UnitOfWork( unit_of_work_code='appreciationengine-20210827-sme-membersvisittotals-v1', licensor=db_licensors['sme'], report=db_reports['membersvisittotals'], report_date='2021-08-27', 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-01T01:01:01.377644+00:00', last_updated_at='2019-11-01T01:00:39.985786+00:00', created_at='2019-11-01T01:01:00.599872+00:00', ) cs = ContentStatus( unit_of_work=uow_1, context='US_Columbia', content_name='US_Columbia.txt', content_status=ContentStatusEnum.MISSING, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) db.session.add_all([uow_1, cs]) db.session.commit() # equal report for i in range(1, 6): uow_2 = UnitOfWork( unit_of_work_code=f'appreciationengine-2021082{i}-sme-membersvisittotals-v1', licensor=db_licensors['sme'], report=db_reports['membersvisittotals'], report_date=f'2021-08-2{i}', version='v1', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.COMPLETE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-01T01:01:01.377644+00:00', last_updated_at='2019-11-01T01:00:39.985786+00:00', created_at='2019-11-01T01:01:00.599872+00:00', ) cs2 = ContentStatus( unit_of_work=uow_2, context='US_Columbia', content_name='US_Columbia.txt', content_status=ContentStatusEnum.COMPLETE, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, meta_data={'last_member_id': f'{i}'} ) cs_other = ContentStatus( unit_of_work=uow_2, context='MX', content_name='MX.txt', content_status=ContentStatusEnum.COMPLETE, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, meta_data={'last_member_id': f'{i*2}'} ) db.session.add_all([uow_2, cs2, cs_other]) db.session.commit() # different report for i in range(1, 10): uow_3 = UnitOfWork( unit_of_work_code=f'appreciationengine-2021082{i}-sme-membersextended-v1', licensor=db_licensors['sme'], report=db_reports['membersextended'], report_date=f'2021-08-2{i}', version='v1', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.COMPLETE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-01T01:01:01.377644+00:00', last_updated_at='2019-11-01T01:00:39.985786+00:00', created_at='2019-11-01T01:01:00.599872+00:00', ) cs3 = ContentStatus( unit_of_work=uow_3, context='US_Columbia', content_name='US_Columbia.txt', content_status=ContentStatusEnum.COMPLETE, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, meta_data={'last_member_id': f'{i*3}'} ) db.session.add_all([uow_3, cs3]) db.session.commit() # cancelled unit uow_3 = UnitOfWork( unit_of_work_code=f'appreciationengine-20210815-sme-membersvisittotals-v1', licensor=db_licensors['sme'], report=db_reports['membersvisittotals'], report_date=f'2021-08-15', version='v1', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.CANCELLED.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-01T01:01:01.377644+00:00', last_updated_at='2019-11-01T01:00:39.985786+00:00', created_at='2019-11-01T01:01:00.599872+00:00', ) cs3 = ContentStatus( unit_of_work=uow_3, context='US_Columbia', content_name='US_Columbia.txt', content_status=ContentStatusEnum.CANCELLED, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, meta_data={'last_member_id': f'555666'} ) db.session.add_all([uow_3, cs3]) db.session.commit() service = ContentStatusService( logger=logger_test, db_conn=db, ) max_id = service.get_max_member_id(cs) assert max_id == 5 @pytest.mark.integration def test_get_max_member_id_filter_date_ok(db, db_licensors, db_reports, logger_test): uow_1 = UnitOfWork( unit_of_work_code='appreciationengine-20210827-sme-membersvisittotals-v1', licensor=db_licensors['sme'], report=db_reports['membersvisittotals'], report_date='2021-08-27', 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-01T01:01:01.377644+00:00', last_updated_at='2019-11-01T01:00:39.985786+00:00', created_at='2019-11-01T01:01:00.599872+00:00', ) cs = ContentStatus( unit_of_work=uow_1, context='US_Columbia', content_name='US_Columbia.txt', content_status=ContentStatusEnum.MISSING, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) db.session.add_all([uow_1, cs]) db.session.commit() # equal report for i in [1, 2, 3, 8, 9]: uow_2 = UnitOfWork( unit_of_work_code=f'appreciationengine-2021082{i}-sme-membersvisittotals-v1', licensor=db_licensors['sme'], report=db_reports['membersvisittotals'], report_date=f'2021-08-2{i}', version='v1', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.COMPLETE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-01T01:01:01.377644+00:00', last_updated_at='2019-11-01T01:00:39.985786+00:00', created_at='2019-11-01T01:01:00.599872+00:00', ) cs2 = ContentStatus( unit_of_work=uow_2, context='US_Columbia', content_name='US_Columbia.txt', content_status=ContentStatusEnum.COMPLETE, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, meta_data={'last_member_id': f'{i}'} ) db.session.add_all([uow_2, cs2]) db.session.commit() service = ContentStatusService( logger=logger_test, db_conn=db, ) max_id = service.get_max_member_id(cs) assert max_id == 3 @pytest.mark.integration def test_get_min_member_id_filter_date_ok(db, db_licensors, db_reports, logger_test): uow_1 = UnitOfWork( unit_of_work_code='appreciationengine-20210827-sme-membersvisittotals-v1', licensor=db_licensors['sme'], report=db_reports['membersvisittotals'], report_date='2021-08-27', 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-01T01:01:01.377644+00:00', last_updated_at='2019-11-01T01:00:39.985786+00:00', created_at='2019-11-01T01:01:00.599872+00:00', ) cs = ContentStatus( unit_of_work=uow_1, context='US_Columbia', content_name='US_Columbia.txt', content_status=ContentStatusEnum.MISSING, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) db.session.add_all([uow_1, cs]) db.session.commit() # equal report for i in [1, 2, 3, 8, 9]: uow_2 = UnitOfWork( unit_of_work_code=f'appreciationengine-2021082{i}-sme-membersvisittotals-v1', licensor=db_licensors['sme'], report=db_reports['membersvisittotals'], report_date=f'2021-08-2{i}', version='v1', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.COMPLETE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-01T01:01:01.377644+00:00', last_updated_at='2019-11-01T01:00:39.985786+00:00', created_at='2019-11-01T01:01:00.599872+00:00', ) cs2 = ContentStatus( unit_of_work=uow_2, context='US_Columbia', content_name='US_Columbia.txt', content_status=ContentStatusEnum.COMPLETE, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, meta_data={ 'last_member_id': f'{i}{i}', 'first_member_id': f'{i}{i}{i}' } ) db.session.add_all([uow_2, cs2]) db.session.commit() service = ContentStatusService( logger=logger_test, db_conn=db, ) max_id = service.get_end_member_id(cs) assert max_id == 888 @pytest.mark.integration def test_get_min_member_id_should_return_none_if_no_data(db, db_licensors, db_reports, logger_test): uow_1 = UnitOfWork( unit_of_work_code='appreciationengine-20210827-sme-membersvisittotals-v1', licensor=db_licensors['sme'], report=db_reports['membersvisittotals'], report_date='2021-08-27', 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-01T01:01:01.377644+00:00', last_updated_at='2019-11-01T01:00:39.985786+00:00', created_at='2019-11-01T01:01:00.599872+00:00', ) cs = ContentStatus( unit_of_work=uow_1, context='US_Columbia', content_name='US_Columbia.txt', content_status=ContentStatusEnum.MISSING, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) db.session.add_all([uow_1, cs]) db.session.commit() service = ContentStatusService( logger=logger_test, db_conn=db, ) max_id = service.get_end_member_id(cs) assert max_id is None @pytest.mark.integration @pytest.mark.parametrize( 'timeslot', [ AEDateTimeRange( lower=datetime(year=2021, month=1, day=1, hour=1), upper=datetime(year=2021, month=1, day=1, hour=6), ).to_dt_range(), None, ] ) def test_get_uow_timeslot(timeslot, db, logger_test, uow_mock): uow_mock.timeslot = timeslot db.session.commit() service = ContentStatusService( logger=logger_test, db_conn=db, ) result = service.get_uow_timeslot(uow_mock.unit_of_work_id) assert result == timeslot @pytest.mark.integration def test_restore_to_complete_status(job_mock, db, db_licensors, db_reports, logger_test): uow_1 = UnitOfWork( unit_of_work_code='appreciationengine-20210827-sme-membersvisittotals-v1', licensor=db_licensors['sme'], report=db_reports['membersvisittotals'], report_date='2021-08-27', 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-01T01:01:01.377644+00:00', last_updated_at='2019-11-01T01:00:39.985786+00:00', created_at='2019-11-01T01:01:00.599872+00:00', ) cs = ContentStatus( unit_of_work=uow_1, context='US_Columbia', content_name='US_Columbia.txt', content_status=ContentStatusEnum.MISSING, content_size=0, failure_count=0, created_at='2019-11-01T01:01:00.599872+00:00', sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) db.session.add_all([uow_1, cs]) db.session.commit() service = ContentStatusService( logger=logger_test, db_conn=db, ) service.restore_to_complete_status(job=job_mock, content_status=cs) assert db.session.query(slz.ContentStatus).count() == 1 cs = db.session.query(slz.ContentStatus).one() assert cs.content_status == ContentStatusEnum.COMPLETE assert cs.latest_job_id == job_mock.job_id