from dataclasses import asdict from datetime import datetime, timezone from unittest.mock import Mock import pytest from db_schema.common import ActivityStatusEnum, CompletenessStatusEnum from db_schema.schemas.slz import ( ContentMetadataStatusEnum, ContentStatus, ContentStatusEnum, DataSource, Report, UnitOfWork, UnitOfWorkPipeline, UnitOfWorkPriorityEnum, ) from slz_job_executor.services.db import DBService @pytest.mark.integration def test_find_units_in_progress_per_dsp(db, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='appreciationengine-20191123-sme-activities-v1', licensor=db_licensors['sme'], report=db_reports['activities'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.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_2 = UnitOfWork( unit_of_work_code='appreciationengine-20191103-sme-activities-v1', licensor=db_licensors['sme'], report=db_reports['activities'], report_date='2019-11-03', version='v1', activity_status=ActivityStatusEnum.IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.MIN_COMPLETE.value, is_force_complete=True, priority=UnitOfWorkPriorityEnum.PRIORITY_8.value, next_run_at='2019-11-01T00:02:01.141397+00:00', created_at='2019-11-01T00:01:00.141397+00:00', last_updated_at='2019-11-01T00:01:55.141397+00:00', ) uow_3 = UnitOfWork( unit_of_work_code='appreciationengine-20191124-sme-activities-v1', licensor=db_licensors['sme'], report=db_reports['activities'], 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.PRIORITY_10.value, next_run_at='2019-11-07T03:12:38.377644+00:00', created_at='2019-11-07T01:11:00.599872+00:00', last_updated_at='2019-11-07T01:11:39.985786+00:00', ) db.session.add_all([uow_1, uow_2, uow_3]) db.session.commit() db_service = DBService( logger=Mock(), db_conn=db, now=Mock(), ) results = db_service.find_units_in_progress_per_dsp() expected = { 'appreciationengine': 2, } assert results == expected @pytest.mark.integration def test_find_scheduled_before(db, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20191123-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', 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_2 = UnitOfWork( unit_of_work_code='spotify-20191103-sme-sub_30_sec_streams-v2', licensor=db_licensors['sme'], report=db_reports['sub_30_sec_streams'], report_date='2019-11-03', version='v2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.MIN_COMPLETE.value, is_force_complete=True, priority=UnitOfWorkPriorityEnum.PRIORITY_8.value, next_run_at='2019-11-01T00:02:01.141397+00:00', created_at='2019-11-01T00:01:00.141397+00:00', last_updated_at='2019-11-01T00:01:55.141397+00:00', ) uow_3 = UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.PRIORITY_10.value, next_run_at='2019-11-07T03:12:38.377644+00:00', created_at='2019-11-07T01:11:00.599872+00:00', last_updated_at='2019-11-07T01:11:39.985786+00:00', ) db.session.add_all([uow_1, uow_2, uow_3]) db.session.commit() db_service = DBService( logger=Mock(), db_conn=db, now=Mock(), ) results = db_service.find_scheduled_before(datetime(year=2019, month=11, day=8, hour=15)) assert [uow_1, uow_2, uow_3] == results results = db_service.find_scheduled_before( datetime(year=2019, month=11, day=8, hour=15), limit=1, ) assert [uow_1] == results results = db_service.find_scheduled_before(datetime(year=2018, month=11, day=8, hour=15)) assert [] == results results = db_service.find_scheduled_before(datetime(year=2019, month=11, day=2, hour=1)) assert [uow_1, uow_2] == results results = db_service.find_scheduled_before( now=datetime(year=2019, month=11, day=8, hour=15), skip_dsp=['apple'] ) assert [uow_2] == results results = db_service.find_scheduled_before( now=datetime(year=2019, month=11, day=8, hour=15), include_only_dsp=['apple'] ) assert [uow_1, uow_3] == results @pytest.mark.integration def test_get_uow_content_status_count(db, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20191123-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', 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', ) db.session.add_all([uow_1]) db.session.commit() db_service = DBService( logger=Mock(), db_conn=db, now=Mock(), ) result = db_service.get_uow_content_status_count(uow_1.unit_of_work_id) assert result == 0 for context in ['AD', 'AE']: cs = ContentStatus( unit_of_work_id=uow_1.unit_of_work_id, context=context, content_name=f'{context}.txt', content_status=ContentStatusEnum.FAILED, content_size=400, 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(cs) db.session.commit() result = db_service.get_uow_content_status_count(uow_1.unit_of_work_id) assert result == 2 uow_3 = UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.PRIORITY_10.value, next_run_at='2019-11-07T03:12:38.377644+00:00', created_at='2019-11-07T01:11:00.599872+00:00', last_updated_at='2019-11-07T01:11:39.985786+00:00', ) result = db_service.get_uow_content_status_count(uow_3.unit_of_work_id) assert result == 0 @pytest.mark.integration def test_get_subtract_contexts(db, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20191123-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', 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', ) db.session.add(uow_1) db.session.commit() expected = { ContentStatusEnum.COMPLETE: ['FR', 'IT', 'RU', 'US'], ContentStatusEnum.CANCELLED: ['BR', 'GB'], ContentStatusEnum.ON_HOLD: ['JP'], } for status, contexts in expected.items(): for context in contexts: cs = ContentStatus( unit_of_work_id=uow_1.unit_of_work_id, context=context, content_name=f'{context}.txt', content_status=status, content_size=400, 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(cs) db.session.commit() db_service = DBService( logger=Mock(), db_conn=db, now=Mock(), ) result = db_service.get_subtract_contexts(unit_of_work_id=uow_1.unit_of_work_id) assert result == expected @pytest.mark.integration def test_create_uow_pipeline(db, logger_test, start_job_uow_mock, test_job_payload): now = datetime(year=2021, month=12, day=24, hour=13, tzinfo=timezone.utc) db.session.add(start_job_uow_mock) db.session.commit() db_service = DBService( logger=logger_test, db_conn=db, now=now, ) assert start_job_uow_mock.unit_of_work_pipeline is None result = db_service.create_or_update_uow_pipeline(start_job_uow_mock, test_job_payload) assert result is True assert start_job_uow_mock.unit_of_work_pipeline is not None assert start_job_uow_mock.unit_of_work_pipeline.config == asdict(test_job_payload) assert start_job_uow_mock.unit_of_work_pipeline.created_at == now @pytest.mark.integration def test_update_uow_pipeline(db, logger_test, start_job_uow_mock, test_job_payload): now = datetime(year=2021, month=12, day=24, hour=13, tzinfo=timezone.utc) now_1 = datetime(year=2021, month=12, day=24, hour=15, tzinfo=timezone.utc) db.session.add(start_job_uow_mock) pipeline = UnitOfWorkPipeline( unit_of_work=start_job_uow_mock, created_at=now, updated_at=now, config={}, ) db.session.add(pipeline) db.session.commit() db_service = DBService( logger=logger_test, db_conn=db, now=now_1, ) assert start_job_uow_mock.unit_of_work_pipeline.config == {} assert start_job_uow_mock.unit_of_work_pipeline.updated_at == now result = db_service.create_or_update_uow_pipeline(start_job_uow_mock, test_job_payload) assert result is True assert start_job_uow_mock.unit_of_work_pipeline.config == asdict(test_job_payload) assert start_job_uow_mock.unit_of_work_pipeline.created_at == now assert start_job_uow_mock.unit_of_work_pipeline.updated_at == now_1