from datetime import date, datetime, timedelta, timezone import pytest from db_schema.schemas import slz from slz_house_keeper.services.db import DBService from slz_house_keeper.services.uow_reset import UowResetService @pytest.mark.integration def test_complete_stuck_should_force_complete_uow_by_active_days_limit( db, logger_test, merged_config_service_test, db_reports, db_licensors ): now = datetime(year=2021, month=1, day=4, hour=12, minute=1, tzinfo=timezone.utc) uow_1 = slz.UnitOfWork( unit_of_work_code='amazonadsupported-20210103-sme-activity-v1', reprocess_id='', licensor=db_licensors['sme'], report=db_reports['activity'], report_date=date(2021, 1, 3), version='v1', 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=now - timedelta(days=181, hours=1, minutes=30), created_at=now - timedelta(days=181, hours=2), last_updated_at=now - timedelta(hours=1), ) db.session.add(uow_1) db.session.commit() db_service = DBService( logger=logger_test, db_conn=db, now=now, ) uow_reset_service = UowResetService( logger=logger_test, now=now, merged_config_service=merged_config_service_test, db_service=db_service, ) assert uow_1.last_updated_at > uow_1.next_run_at db.session.query(slz.UnitOfWork ).filter(slz.ContentStatus.unit_of_work_id == uow_1.unit_of_work_id).update( {'activity_status': slz.ActivityStatusEnum.NOT_IN_PROGRESS.value}, synchronize_session=False ) db.session.commit() result = uow_reset_service.complete_stuck_in_active_uows() assert result == [uow_1.readable] assert uow_1.completeness_status == slz.CompletenessStatusEnum.COMPLETE assert uow_1.last_updated_at == now assert uow_1.is_force_complete is True @pytest.mark.integration def test_complete_stuck_should_not_force_complete_uow_by_active_days_limit( db, logger_test, merged_config_service_test, db_reports, db_licensors ): now = datetime(year=2021, month=1, day=4, hour=12, minute=1, tzinfo=timezone.utc) uow_1 = slz.UnitOfWork( unit_of_work_code='amazonadsupported-20210103-sme-activity-v123', reprocess_id='', licensor=db_licensors['sme'], report=db_reports['activity'], report_date=date(2021, 1, 3), version='v1', activity_status=slz.ActivityStatusEnum.IN_PROGRESS.value, completeness_status=slz.CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=slz.UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at=now - timedelta(days=181, hours=1, minutes=30), created_at=now - timedelta(days=181, hours=2), last_updated_at=now - timedelta(hours=1), ) db.session.add(uow_1) db.session.commit() db_service = DBService( logger=logger_test, db_conn=db, now=now, ) uow_reset_service = UowResetService( logger=logger_test, now=now, merged_config_service=merged_config_service_test, db_service=db_service, ) assert uow_1.last_updated_at > uow_1.next_run_at result = uow_reset_service.complete_stuck_in_active_uows() assert result == [] assert uow_1.completeness_status == slz.CompletenessStatusEnum.ACTIVE assert uow_1.is_force_complete is False @pytest.mark.integration def test_complete_stuck_should_not_schedule_if_config_not_found( db, logger_test, merged_config_service_test, db_reports, db_licensors ): now = datetime(year=2021, month=1, day=4, hour=12, minute=1, tzinfo=timezone.utc) uow_1 = slz.UnitOfWork( unit_of_work_code='amazonadsupported-20210103-sme-activity-v2222', reprocess_id='', licensor=db_licensors['sme'], report=db_reports['activity'], report_date=date(2021, 1, 3), version='v2222', 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=now - timedelta(hours=1, minutes=30), created_at=now - timedelta(hours=2), last_updated_at=now - timedelta(hours=1), ) db.session.add(uow_1) db.session.commit() db_service = DBService( logger=logger_test, db_conn=db, now=now, ) uow_reset_service = UowResetService( logger=logger_test, now=now, merged_config_service=merged_config_service_test, db_service=db_service, ) result = uow_reset_service.complete_stuck_in_active_uows() assert result == [] @pytest.mark.integration @pytest.mark.parametrize( 'content_statuses, expected_is_completed', [ ({ 'US': slz.ContentStatusEnum.COMPLETE, 'AT': slz.ContentStatusEnum.COMPLETE }, True), ({ 'US': slz.ContentStatusEnum.COMPLETE, 'AT': slz.ContentStatusEnum.MISSING }, False), ({ 'US': slz.ContentStatusEnum.COMPLETE, 'AT': slz.ContentStatusEnum.FAILED }, False), ({ 'US': slz.ContentStatusEnum.COMPLETE, 'AT': slz.ContentStatusEnum.ON_HOLD }, False), ] ) def test_complete_stuck_by_content_statuses( content_statuses, expected_is_completed, db, logger_test, merged_config_service_test, db_reports, db_licensors, ): now = datetime(year=2021, month=1, day=4, hour=12, minute=1, tzinfo=timezone.utc) uow_1 = slz.UnitOfWork( unit_of_work_code='amazonadsupported-20210103-sme-activity-v1', reprocess_id='', licensor=db_licensors['sme'], report=db_reports['activity'], report_date=date(2021, 1, 3), version='v1', 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=now - timedelta(days=1, hours=1, minutes=30), created_at=now - timedelta(days=1, hours=2), last_updated_at=now - timedelta(hours=1), ) db.session.add(uow_1) for context, status in content_statuses.items(): cs = slz.ContentStatus( unit_of_work=uow_1, context=context, content_name='us.txt', content_status=status, created_at=now, failure_count=0, sub_content='{}', metadata_process_status=slz.ContentMetadataStatusEnum.NOT_QUEUED.value, ) db.session.add(cs) db.session.commit() db_service = DBService( logger=logger_test, db_conn=db, now=now, ) uow_reset_service = UowResetService( logger=logger_test, now=now, merged_config_service=merged_config_service_test, db_service=db_service, ) result = uow_reset_service.complete_stuck_in_active_uows() assert bool(result) is expected_is_completed assert (uow_1.completeness_status == slz.CompletenessStatusEnum.COMPLETE) == \ expected_is_completed assert (uow_1.last_updated_at == now) == expected_is_completed @pytest.mark.integration @pytest.mark.parametrize( 'uow_completeness_status, expected_group_status', [ (slz.CompletenessStatusEnum.COMPLETE, slz.CompletenessStatusEnum.COMPLETE), (slz.CompletenessStatusEnum.ACTIVE, slz.CompletenessStatusEnum.ACTIVE), (slz.CompletenessStatusEnum.MIN_COMPLETE, slz.CompletenessStatusEnum.ACTIVE), (slz.CompletenessStatusEnum.CANCELLED, slz.CompletenessStatusEnum.COMPLETE), ] ) def test_complete_uow_group( uow_completeness_status, expected_group_status, db, logger_test, now, merged_config_service_test, unit_of_work_stubs, ): uow_1 = slz.UnitOfWork(**unit_of_work_stubs[1]) uow_1.completeness_status = uow_completeness_status uow_2 = slz.UnitOfWork(**unit_of_work_stubs[2]) uow_3 = slz.UnitOfWork(**unit_of_work_stubs[2]) uow_3.reprocess_id = '1234' uow_3.completeness_status = slz.CompletenessStatusEnum.MIN_COMPLETE db.session.add_all([uow_1, uow_2, uow_3]) group_1 = slz.UnitOfWorkGroup( created_at=now, completeness_status=slz.CompletenessStatusEnum.ACTIVE.value ) group_2 = slz.UnitOfWorkGroup( created_at=now, completeness_status=slz.CompletenessStatusEnum.ACTIVE.value ) db.session.add_all([group_1, group_2]) db.session.commit() pipeline_1 = slz.UnitOfWorkPipeline( unit_of_work_id=uow_1.unit_of_work_id, unit_of_work_group_id=group_1.unit_of_work_group_id, created_at=now, config={}, ) pipeline_2 = slz.UnitOfWorkPipeline( unit_of_work_id=uow_2.unit_of_work_id, unit_of_work_group_id=group_1.unit_of_work_group_id, created_at=now, config={}, ) pipeline_3 = slz.UnitOfWorkPipeline( unit_of_work_id=uow_3.unit_of_work_id, unit_of_work_group_id=group_2.unit_of_work_group_id, created_at=now, config={}, ) db.session.add_all([pipeline_1, pipeline_2, pipeline_3]) db.session.commit() db_service = DBService( logger=logger_test, db_conn=db, now=now, ) uow_reset_service = UowResetService( logger=logger_test, now=now, merged_config_service=merged_config_service_test, db_service=db_service, ) uow_reset_service.complete_stuck_in_active_uow_group() assert group_1.completeness_status == expected_group_status assert group_2.completeness_status == slz.CompletenessStatusEnum.ACTIVE @pytest.mark.integration @pytest.mark.parametrize( 'timedelta_hours, expected_is_completed', [ (2, False), (6, False), (24, True), ] ) def test_complete_by_uow_complete_after_hours( timedelta_hours, expected_is_completed, db, logger_test, merged_config_service_test, db_reports, db_licensors, ): now = datetime(year=2021, month=12, day=10, hour=12, minute=1, tzinfo=timezone.utc) uow_1 = slz.UnitOfWork( unit_of_work_code='appreciationengine-20211210-sme-request_to_forget-v1', reprocess_id='', licensor=db_licensors['sme'], report=db_reports['request_to_forget'], report_date=date(2021, 12, 10), version='v1', 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=now - timedelta(hours=1, minutes=30), created_at=now - timedelta(hours=2), last_updated_at=now - timedelta(hours=1), ) db.session.add(uow_1) db.session.commit() db_service = DBService( logger=logger_test, db_conn=db, now=now, ) uow_reset_service = UowResetService( logger=logger_test, now=now + timedelta(hours=timedelta_hours), merged_config_service=merged_config_service_test, db_service=db_service, ) assert uow_1.completeness_status != slz.CompletenessStatusEnum.COMPLETE uow_reset_service.complete_stuck_in_active_uows() assert (uow_1.completeness_status == slz.CompletenessStatusEnum.COMPLETE) == \ expected_is_completed assert ( uow_1.last_updated_at == now + timedelta(hours=timedelta_hours) ) == expected_is_completed