# pylint: disable=unused-argument from datetime import date, datetime, timezone import pytest from db_schema.schemas import slz from slz_job_manager.services.db import CommandDBService, QueryDBService @pytest.mark.integration def test_create_uow_without_timeframe(db, clean_db, unit_of_work_stubs, logger_test): now = datetime(year=2021, month=1, day=1, tzinfo=timezone.utc) db_service = CommandDBService( logger=logger_test, db_conn=db, now=now, ) uow_id = 'apple-20191124-theorchard-amEvent-v1_2' # no unit in db count = db.session.query(slz.UnitOfWork).filter(slz.UnitOfWork.unit_of_work_code == uow_id ).count() assert count == 0 # create uow result = db_service.create_uow(uow_id, now, now, now) assert result is True # check that it was successfully created count = db.session.query(slz.UnitOfWork).filter(slz.UnitOfWork.unit_of_work_code == uow_id ).count() assert count == 1 # try to create the same uow one more time result = db_service.create_uow(uow_id, now, now, now) assert result is False # same uow should not be created count = db.session.query(slz.UnitOfWork).filter(slz.UnitOfWork.unit_of_work_code == uow_id ).count() assert count == 1 @pytest.mark.integration def test_create_uow_with_timeframe(db, clean_db, unit_of_work_stubs, logger_test): now = datetime(year=2021, month=1, day=1, tzinfo=timezone.utc) db_service = CommandDBService( logger=logger_test, db_conn=db, now=now, ) uow_id = 'apple-20191124-theorchard-amEvent-v1_2' group = db_service.create_uow_group(uow_id, now) # no unit in db count = db.session.query(slz.UnitOfWork).filter(slz.UnitOfWork.unit_of_work_code == uow_id ).count() assert count == 0 # create uow result = db_service.create_uow( uow_id=uow_id, created_at=now, last_updated_at=now, next_run_at=now, group=group, timeframe=[ datetime(year=2021, month=1, day=1, hour=1), datetime(year=2021, month=1, day=1, hour=4), ] ) assert result is True # check that it was successfully created count = db.session.query(slz.UnitOfWork).filter(slz.UnitOfWork.unit_of_work_code == uow_id ).count() assert count == 1 # try to create the same uow one more time result = db_service.create_uow( uow_id=uow_id, created_at=now, last_updated_at=now, next_run_at=now, group=group, timeframe=[ datetime(year=2021, month=1, day=1, hour=1), datetime(year=2021, month=1, day=1, hour=4), ] ) assert result is False # same uow should not be created count = db.session.query(slz.UnitOfWork).filter(slz.UnitOfWork.unit_of_work_code == uow_id ).count() assert count == 1 @pytest.mark.integration def test_create_group(db, clean_db, logger_test): now = datetime(year=2021, month=1, day=1, tzinfo=timezone.utc) db_service = CommandDBService( logger=logger_test, db_conn=db, now=now, ) count = db.session.query(slz.UnitOfWorkGroup).count() assert count == 0 group = db_service.create_uow_group('apple-20191124-theorchard-amEvent-v1_2', now) assert isinstance(group, slz.UnitOfWorkGroup) count = db.session.query(slz.UnitOfWorkGroup).count() assert count == 1 @pytest.mark.integration def test_find_uow_incomplete_to_reschedule( db, clean_db, db_ro, unit_of_work_stubs, logger_test, db_licensors, db_reports ): db_service = QueryDBService( logger=logger_test, db_conn=db_ro, ) uow_1 = slz.UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date=date(year=2019, month=11, day=24), version='v1_2', activity_status=slz.ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=slz.CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=slz.UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at=datetime(year=2019, month=11, day=7, hour=1, minute=10, tzinfo=timezone.utc), created_at=datetime(year=2019, month=11, day=7, hour=1, minute=1, tzinfo=timezone.utc), last_updated_at=datetime( year=2019, month=11, day=21, hour=18, minute=41, tzinfo=timezone.utc ), ) db.session.add(uow_1) uow_2 = slz.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=date(year=2019, month=11, day=3), version='v2', activity_status=slz.ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=slz.CompletenessStatusEnum.MIN_COMPLETE.value, is_force_complete=True, priority=slz.UnitOfWorkPriorityEnum.PRIORITY_8.value, next_run_at=datetime(year=2019, month=11, day=4, hour=2, minute=55, tzinfo=timezone.utc), created_at=datetime(year=2019, month=11, day=4, hour=2, minute=55, tzinfo=timezone.utc), last_updated_at=datetime( year=2019, month=11, day=4, hour=2, minute=56, tzinfo=timezone.utc ), ) db.session.add(uow_2) uow_3 = slz.UnitOfWork( unit_of_work_code='spotify-20191104-sme-sub_30_sec_streams-v2', licensor=db_licensors['sme'], report=db_reports['sub_30_sec_streams'], report_date=date(year=2019, month=11, day=4), version='v2', activity_status=slz.ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=slz.CompletenessStatusEnum.ON_HOLD.value, is_force_complete=True, priority=slz.UnitOfWorkPriorityEnum.PRIORITY_8.value, next_run_at=datetime(year=2019, month=11, day=5, hour=2, minute=55, tzinfo=timezone.utc), created_at=datetime(year=2019, month=11, day=4, hour=2, minute=55, tzinfo=timezone.utc), last_updated_at=datetime( year=2019, month=11, day=4, hour=2, minute=55, tzinfo=timezone.utc ), ) db.session.add(uow_3) db.session.commit() results = db_service.find_uow_incomplete_to_reschedule() assert results == [(uow.unit_of_work_id, uow.unit_of_work_code) for uow in (uow_1, uow_2)] @pytest.mark.integration def test_update(db, clean_db, unit_of_work_stubs, logger_test): db_service = CommandDBService( logger=logger_test, db_conn=db, now=datetime(year=2021, month=1, day=1, tzinfo=timezone.utc) ) uow_1 = slz.UnitOfWork(**unit_of_work_stubs[1]) db.session.add(uow_1) db.session.commit() assert uow_1.version != 'v1_3' db_service.update( uow_1, activity_status=slz.ActivityStatusEnum.IN_PROGRESS.value, version='v1_3', ) found = db.session.query(slz.UnitOfWork ).filter(slz.UnitOfWork.unit_of_work_id == uow_1.unit_of_work_id).one() assert found.activity_status == slz.ActivityStatusEnum.IN_PROGRESS assert found.version == 'v1_3' @pytest.mark.integration def test_create_reprocessing_uow(db, clean_db, unit_of_work_stubs, logger_test): now = datetime(year=2021, month=1, day=1, tzinfo=timezone.utc) db_service = CommandDBService( logger=logger_test, db_conn=db, now=now, ) reprocessing_id = '20210101T101315' uow_1 = slz.UnitOfWork(**unit_of_work_stubs[1]) db.session.add(uow_1) db.session.commit() # no unit in db assert db.session.query(slz.UnitOfWork).count() == 1 # create uow result = db_service.create_reprocessing_uow( uow_1, reprocessing_id, slz.UnitOfWorkPriorityEnum.PRIORITY_6, ) assert result is not None # check that it was successfully created assert db.session.query(slz.UnitOfWork).count() == 2 # try to create uow with the same reprocess id result = db_service.create_reprocessing_uow( uow_1, reprocessing_id, slz.UnitOfWorkPriorityEnum.PRIORITY_6, ) assert result is None # uow should not be created assert db.session.query(slz.UnitOfWork).count() == 2