# pylint: disable=unused-argument from datetime import datetime, timezone from unittest.mock import Mock import pytest from db_schema.schemas import slz from slz_job_manager.exceptions import InvalidPayloadError from slz_job_manager.services.db import CommandDBService, QueryDBService from slz_job_manager.services.uow_reprocess import UowReprocessService @pytest.mark.integration def test_reprocess_without_group_success(db, clean_db, db_ro, unit_of_work_stubs, logger_test): now = datetime(year=2021, month=1, day=5, tzinfo=timezone.utc) cmd_db_service = CommandDBService( logger=logger_test, db_conn=db, now=now, ) query_db_service = QueryDBService( logger=logger_test, db_conn=db_ro, ) uow_1 = slz.UnitOfWork(**unit_of_work_stubs[1]) db.session.add(uow_1) db.session.commit() sns_client = Mock() uow_reprocessor = UowReprocessService( logger=logger_test, now=now, command_db_service=cmd_db_service, query_db_service=query_db_service, sns_client=sns_client, sns_receiver='test_topic', ) payload = [{ 'unit_of_work_id': uow_1.unit_of_work_id, 'priority': 8, }] assert db.session.query(slz.UnitOfWork).count() == 1 result = uow_reprocessor.reprocess_uow(payload) assert result == ['apple-20191124-theorchard-amEvent-v1_2-rp20210105T000000'] assert db.session.query(slz.UnitOfWork).count() == 2 assert db.session.query(slz.UnitOfWorkGroup).count() == 0 assert db.session.query(slz.UnitOfWorkPipeline).count() == 0 @pytest.mark.integration def test_reprocess_with_group_and_pipeline_success( db, clean_db, db_ro, unit_of_work_stubs, logger_test ): now = datetime(year=2021, month=1, day=5, tzinfo=timezone.utc) cmd_db_service = CommandDBService( logger=logger_test, db_conn=db, now=now, ) query_db_service = QueryDBService( logger=logger_test, db_conn=db_ro, ) uow_1 = slz.UnitOfWork(**unit_of_work_stubs[1]) db.session.add(uow_1) db.session.commit() group = cmd_db_service.create_uow_group(uow_1.readable, now) group.completeness_status = slz.CompletenessStatusEnum.COMPLETE db.session.commit() pipeline = cmd_db_service.create_uow_pipeline(uow_1, group) sns_client = Mock() uow_reprocessor = UowReprocessService( logger=logger_test, now=now, command_db_service=cmd_db_service, query_db_service=query_db_service, sns_client=sns_client, sns_receiver='test_topic', ) assert db.session.query(slz.UnitOfWork).count() == 1 result = uow_reprocessor.reprocess_uow([{ 'unit_of_work_id': uow_1.unit_of_work_id, 'priority': 8, }]) assert result == ['apple-20191124-theorchard-amEvent-v1_2-rp20210105T000000'] assert db.session.query(slz.UnitOfWork).count() == 2 new: slz.UnitOfWork = db.session.query( slz.UnitOfWork ).filter(slz.UnitOfWork.reprocess_id == '20210105T000000').one() assert new.unit_of_work_pipeline.unit_of_work_group == group assert new.unit_of_work_pipeline != pipeline assert new.unit_of_work_pipeline.unit_of_work_group.completeness_status == \ slz.CompletenessStatusEnum.ACTIVE @pytest.mark.integration def test_reprocess_no_uow_found_in_db(db, clean_db, db_ro, unit_of_work_stubs, logger_test): now = datetime(year=2021, month=1, day=5, tzinfo=timezone.utc) cmd_db_service = CommandDBService( logger=logger_test, db_conn=db, now=now, ) query_db_service = QueryDBService( logger=logger_test, db_conn=db_ro, ) uow_1 = slz.UnitOfWork(**unit_of_work_stubs[1]) db.session.add(uow_1) db.session.commit() sns_client = Mock() uow_reprocessor = UowReprocessService( logger=logger_test, now=now, command_db_service=cmd_db_service, query_db_service=query_db_service, sns_client=sns_client, sns_receiver='test_topic', ) payload = [{'unit_of_work_id': 123}] assert db.session.query(slz.UnitOfWork).count() == 1 result = uow_reprocessor.reprocess_uow(payload) assert result == [] # no new uow were created assert db.session.query(slz.UnitOfWork).count() == 1 @pytest.mark.integration def test_reprocess_bad_payload(db, clean_db, db_ro, unit_of_work_stubs, logger_test): now = datetime(year=2021, month=1, day=5, tzinfo=timezone.utc) cmd_db_service = CommandDBService( logger=logger_test, db_conn=db, now=now, ) query_db_service = QueryDBService( logger=logger_test, db_conn=db_ro, ) uow_1 = slz.UnitOfWork(**unit_of_work_stubs[1]) db.session.add(uow_1) db.session.commit() sns_client = Mock() uow_reprocessor = UowReprocessService( logger=logger_test, now=now, command_db_service=cmd_db_service, query_db_service=query_db_service, sns_client=sns_client, sns_receiver='test_topic', ) payload = [{'uow_id_bad': 'apple-20191124-sme-amEvent-v1_3'}] with pytest.raises(InvalidPayloadError): uow_reprocessor.reprocess_uow(payload) @pytest.mark.integration def test_reprocess_bad_priority_in_payload(db, clean_db, db_ro, unit_of_work_stubs, logger_test): now = datetime(year=2021, month=1, day=5, tzinfo=timezone.utc) cmd_db_service = CommandDBService( logger=logger_test, db_conn=db, now=now, ) query_db_service = QueryDBService( logger=logger_test, db_conn=db_ro, ) uow_1 = slz.UnitOfWork(**unit_of_work_stubs[1]) db.session.add(uow_1) db.session.commit() sns_client = Mock() uow_reprocessor = UowReprocessService( logger=logger_test, now=now, command_db_service=cmd_db_service, query_db_service=query_db_service, sns_client=sns_client, sns_receiver='test_topic', ) payload = [{ 'unit_of_work_id': uow_1.unit_of_work_id, 'priority': 25, }] with pytest.raises(ValueError): uow_reprocessor.reprocess_uow(payload)