from datetime import datetime from unittest.mock import Mock import pytest from db_schema.common import ActivityStatusEnum, CompletenessStatusEnum from db_schema.schemas.slz import ( ContentStatusEnum, DataSource, Report, UnitOfWork, UnitOfWorkPriorityEnum, ) from slz_job_executor.services.job_preparer import JobPreparerService @pytest.fixture def preparer_uow_mock(): uow = UnitOfWork( unit_of_work_code='amazonadsupported-20191124-sme-activity-v1', reprocess_id='', licensor_id=1, report_id=1, report_date=datetime(2019, 11, 24).date(), 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', ) uow.report = Report( report_id=1, report_name='activity', is_active=True, data_source=DataSource( data_source_id=1, data_source_name='amazonadsupported', is_active=True, ) ) return uow @pytest.mark.parametrize( 'get_uow_content_status_count_value, get_subtract_contexts_value, expected_contexts, ' 'reprocess_id, expected_is_reprocess_flag', [ # newly created unit (0, {}, ['US', 'AT'], '', False), # newly created reprocess unit (0, {}, ['US', 'AT'], '20210202123456', True), # not new unit, but no content statuses in excluding statuses (2, {}, ['US', 'AT'], '', False), # not new unit, some content statuses are completed (2, { ContentStatusEnum.COMPLETE: ['US'] }, ['AT'], '', False), # not new reprocess unit, some content statuses are completed (2, { ContentStatusEnum.COMPLETE: ['US'] }, ['AT'], '20210202123456', False), ] ) def test_job_preparer_start_job( get_uow_content_status_count_value, get_subtract_contexts_value, expected_contexts, reprocess_id, expected_is_reprocess_flag, logger_test, merged_config_service_test, preparer_uow_mock, ): preparer_uow_mock.reprocess_id = reprocess_id now = datetime(2019, 11, 24) db_service = Mock() job_starter = Mock() job_preparer = JobPreparerService( logger=logger_test, db_service=db_service, merged_config_service=merged_config_service_test, job_starter_service=job_starter, ) db_service.find_scheduled_before.return_value = [preparer_uow_mock] db_service.get_uow_content_status_count.return_value = get_uow_content_status_count_value db_service.get_subtract_contexts.return_value = get_subtract_contexts_value db_service.find_units_in_progress_per_dsp.return_value = {} db_service.get_active_daily_uow_per_dsp.return_value = [] job_preparer.start_jobs(now) job_starter.start_job.assert_called_with( unit_of_work=preparer_uow_mock, config=merged_config_service_test.get(preparer_uow_mock.readable), contexts=expected_contexts, is_reprocessing=expected_is_reprocess_flag, ) def test_job_preparer_no_config_no_job_started( logger_test, merged_config_service_test, preparer_uow_mock ): now = datetime(2019, 11, 24) db_service = Mock() job_starter = Mock() job_preparer = JobPreparerService( logger=logger_test, db_service=db_service, merged_config_service=merged_config_service_test, job_starter_service=job_starter, ) preparer_uow_mock.unit_of_work_code = 'unknown' db_service.find_scheduled_before.return_value = [preparer_uow_mock] db_service.get_active_daily_uow_per_dsp.return_value = [] job_preparer.start_jobs(now) job_starter.start_job.assert_not_called() def test_start_job_non_slz_flow(logger_test, merged_config_service_test, preparer_uow_mock): now = datetime(2019, 11, 24) db_service = Mock() job_starter = Mock() job_preparer = JobPreparerService( logger=logger_test, db_service=db_service, merged_config_service=merged_config_service_test, job_starter_service=job_starter, ) preparer_uow_mock.unit_of_work_code = 'appreciationengine-20191124-sme-request_to_forget-v1' db_service.find_scheduled_before.return_value = [preparer_uow_mock] db_service.get_uow_content_status_count.return_value = 0 db_service.get_subtract_contexts.return_value = {} db_service.get_active_daily_uow_per_dsp.return_value = [] job_preparer.start_jobs(now) job_starter.start_job.assert_called_with( unit_of_work=preparer_uow_mock, config=merged_config_service_test.get(preparer_uow_mock.readable), # no complete criteria in config. This should be empty contexts=[] ) def test_get_dsp_concurrency_limits(logger_test, merged_config_service_test, preparer_uow_mock): job_preparer = JobPreparerService( logger=logger_test, db_service=Mock(), merged_config_service=merged_config_service_test, job_starter_service=Mock(), ) results = job_preparer.get_dsp_concurrency_limits() expected = {'appreciationengine': 5, 'brandwatch': 1} assert results == expected