# pylint: disable=unused-argument from datetime import date, datetime, timedelta, timezone import pytest from db_schema.schemas.admin import AuditLog, AuditLogTypes from db_schema.schemas.apps import ActivityStatusEnum as AppsActivityStatusEnum from db_schema.schemas.apps import CompletenessStatusEnum as AppsCompletenessStatusEnum from db_schema.schemas.apps import DBXJobStatusEnum from db_schema.schemas.apps import UnitOfWork as AppsUnitOfWork from db_schema.schemas.apps import UnitOfWorkTypeEnum as AppsUnitOfWorkTypeEnum from delphi_slz_admin.modelviews.apps_unit_of_work import ( AppsUnitOfWorkModelView, FilterActivityStatus, FilterCompletenessStatus, FilterDbxExecutionIdIsNone, FilterLatestJobIdIsNone, FilterLatestJobIdIsNotNone, FilterPriority, FilterReportDateEarlierThan, FilterReportDateLaterThan, FilterUnitOfWorkId, FilterUnitOfWorkIdReversed, FilterUnitOfWorkType, ) FREEZED_DT = datetime(2015, 10, 15, 12, 23, 45, tzinfo=timezone.utc) @pytest.mark.freeze_time(FREEZED_DT) def test_action_restart_failed_backfill__ok(db, clean_db, app, apps_unit_of_work): db.begin() modelview = AppsUnitOfWorkModelView(db) modelview.action_restart_failed_backfill([apps_unit_of_work.unit_of_work_id]) db.rollback() db.refresh(apps_unit_of_work) assert apps_unit_of_work.completeness_status == AppsCompletenessStatusEnum.QUEUED assert apps_unit_of_work.last_updated_at == FREEZED_DT assert apps_unit_of_work.next_run_at == FREEZED_DT assert apps_unit_of_work.latest_job_id == ( f'{apps_unit_of_work.unit_of_work_code}_{apps_unit_of_work.unit_of_work_id}_' f'{FREEZED_DT.strftime("%Y%m%dT%H.%M.%S")}' ) @pytest.mark.freeze_time(FREEZED_DT) @pytest.mark.parametrize( 'status', [x for x in AppsCompletenessStatusEnum if x != AppsCompletenessStatusEnum.FAILED] ) def test_action_restart_failed_backfill__invalid_completeness_status( db, clean_db, app, apps_unit_of_work, status, mocker ): patched_flash = mocker.patch( 'delphi_slz_admin.modelviews.apps_unit_of_work.flash', spec=lambda *args, **kwargs: None ) apps_unit_of_work.completeness_status = status db.begin() modelview = AppsUnitOfWorkModelView(db) modelview.action_restart_failed_backfill([apps_unit_of_work.unit_of_work_id]) db.rollback() db.refresh(apps_unit_of_work) assert apps_unit_of_work.completeness_status == status assert apps_unit_of_work.last_updated_at != FREEZED_DT assert apps_unit_of_work.next_run_at != FREEZED_DT patched_flash.assert_called_once_with( f'Some items are not supposed to be restarted. Completeness Status should be ' f'equal to {AppsCompletenessStatusEnum.FAILED.value}. ' f'UoW Type should be equal to {AppsUnitOfWorkTypeEnum.BACKFILL.value}.', category='error' ) @pytest.mark.freeze_time(FREEZED_DT) @pytest.mark.parametrize( 'uow_type', [x for x in AppsUnitOfWorkTypeEnum if x != AppsUnitOfWorkTypeEnum.BACKFILL] ) def test_action_restart_failed_backfill__invalid_uow_type( db, clean_db, app, apps_unit_of_work, uow_type, mocker ): patched_flash = mocker.patch( 'delphi_slz_admin.modelviews.apps_unit_of_work.flash', spec=lambda *args, **kwargs: None ) apps_unit_of_work.unit_of_work_type = uow_type db.begin() modelview = AppsUnitOfWorkModelView(db) modelview.action_restart_failed_backfill([apps_unit_of_work.unit_of_work_id]) db.rollback() db.refresh(apps_unit_of_work) assert apps_unit_of_work.unit_of_work_type == uow_type assert apps_unit_of_work.last_updated_at != FREEZED_DT assert apps_unit_of_work.next_run_at != FREEZED_DT patched_flash.assert_called_once_with( f'Some items are not supposed to be restarted. Completeness Status should be ' f'equal to {AppsCompletenessStatusEnum.FAILED.value}. ' f'UoW Type should be equal to {AppsUnitOfWorkTypeEnum.BACKFILL.value}.', category='error' ) @pytest.mark.freeze_time(FREEZED_DT) def test_action_restart_failed_backfill__audit_log_created(db, clean_db, app, apps_unit_of_work): db.begin() modelview = AppsUnitOfWorkModelView(db) modelview.action_restart_failed_backfill([apps_unit_of_work.unit_of_work_id]) db.rollback() audit_log: AuditLog = db.query(AuditLog).first() assert audit_log.event_type == AuditLogTypes.APPS_UOW_RESTART_FAILED_BACKFILL assert audit_log.created_at == FREEZED_DT assert audit_log.extra_data == {'apps_unit_of_work_ids': [apps_unit_of_work.unit_of_work_id]} def test_action_restart_failed_backfill__latest_job_id_truncated( db, clean_db, app, apps_unit_of_work ): apps_unit_of_work.unit_of_work_code = 'a very long code' * 40 db.begin() modelview = AppsUnitOfWorkModelView(db) modelview.action_restart_failed_backfill([apps_unit_of_work.unit_of_work_id]) db.rollback() db.refresh(apps_unit_of_work) assert len(apps_unit_of_work.latest_job_id) == AppsUnitOfWorkModelView.truncate_limit def test_filter_unit_of_work_id__apply_included(db, clean_db, apps_unit_of_work): filter_ = FilterUnitOfWorkId() initial_query = db.query(AppsUnitOfWork) cleaned_value = filter_.clean( f'{apps_unit_of_work.unit_of_work_id}, {apps_unit_of_work.unit_of_work_id + 1}' ) filtered_query = filter_.apply(initial_query, value=cleaned_value) assert initial_query.count() == 1 assert filtered_query.count() == 1 def test_filter_unit_of_work_id__apply_excluded(db, clean_db, apps_unit_of_work): filter_ = FilterUnitOfWorkId() initial_query = db.query(AppsUnitOfWork) cleaned_value = filter_.clean( f'{apps_unit_of_work.unit_of_work_id + 1}, {apps_unit_of_work.unit_of_work_id + 2}' ) filtered_query = filter_.apply(initial_query, value=cleaned_value) assert initial_query.count() == 1 assert filtered_query.count() == 0 def test_filter_unit_of_work_id_reversed_apply_included(db, clean_db, apps_unit_of_work): filter_ = FilterUnitOfWorkIdReversed() initial_query = db.query(AppsUnitOfWork) cleaned_value = filter_.clean( f'{apps_unit_of_work.unit_of_work_id}, {apps_unit_of_work.unit_of_work_id + 1}' ) filtered_query = filter_.apply(initial_query, value=cleaned_value) assert initial_query.count() == 1 assert filtered_query.count() == 0 def test_filter_unit_of_work_id_reversed__apply_excluded(db, clean_db, apps_unit_of_work): filter_ = FilterUnitOfWorkIdReversed() initial_query = db.query(AppsUnitOfWork) cleaned_value = filter_.clean( f'{apps_unit_of_work.unit_of_work_id + 1}, {apps_unit_of_work.unit_of_work_id + 2}' ) filtered_query = filter_.apply(initial_query, value=cleaned_value) assert initial_query.count() == 1 assert filtered_query.count() == 1 @pytest.mark.parametrize('value', ['foo', 'foo, 123']) def test_filter_unit_of_work_id__invalid_values(db, clean_db, value): filter_ = FilterUnitOfWorkId() with pytest.raises(ValueError): filter_.clean(value) def test_filter_unit_of_work_type__apply_included(db, clean_db, apps_unit_of_work): filter_ = FilterUnitOfWorkType() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=AppsUnitOfWorkTypeEnum.BACKFILL.value) assert initial_query.count() == 1 assert filtered_query.count() == 1 def test_filter_unit_of_work_type__apply_excluded(db, clean_db, apps_unit_of_work): filter_ = FilterUnitOfWorkType() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=AppsUnitOfWorkTypeEnum.DAILY.value) assert initial_query.count() == 1 assert filtered_query.count() == 0 def test_filter_completeness_status__apply_included(db, clean_db, apps_unit_of_work): filter_ = FilterCompletenessStatus() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=AppsCompletenessStatusEnum.FAILED.value) assert initial_query.count() == 1 assert filtered_query.count() == 1 def test_filter_completeness_status__apply_excluded(db, clean_db, apps_unit_of_work): filter_ = FilterCompletenessStatus() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=AppsCompletenessStatusEnum.COMPLETE.value) assert initial_query.count() == 1 assert filtered_query.count() == 0 def test_filter_activity_status__apply_included(db, clean_db, apps_unit_of_work): filter_ = FilterActivityStatus() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply( initial_query, value=AppsActivityStatusEnum.NOT_IN_PROGRESS.value ) assert initial_query.count() == 1 assert filtered_query.count() == 1 def test_filter_activity_status__apply_excluded(db, clean_db, apps_unit_of_work): filter_ = FilterActivityStatus() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=AppsActivityStatusEnum.IN_PROGRESS.value) assert initial_query.count() == 1 assert filtered_query.count() == 0 def test_filter_priority__apply_included(db, clean_db, apps_unit_of_work): filter_ = FilterPriority() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value='10') assert initial_query.count() == 1 assert filtered_query.count() == 1 def test_filter_priority__apply_excluded(db, clean_db, apps_unit_of_work): filter_ = FilterPriority() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value='5') assert initial_query.count() == 1 assert filtered_query.count() == 0 @pytest.mark.parametrize('value', ['foo', 'foo, 123']) def test_filter_priority__invalid_values(db, clean_db, value): filter_ = FilterPriority() with pytest.raises(ValueError): filter_.clean(value) def test_filter_latest_job_id_is_none__apply_included(db, clean_db, apps_unit_of_work): apps_unit_of_work.latest_job_id = None filter_ = FilterLatestJobIdIsNone() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=None) assert initial_query.count() == 1 assert filtered_query.count() == 1 def test_filter_latest_job_id_is_none__apply_excluded(db, clean_db, apps_unit_of_work): apps_unit_of_work.latest_job_id = 'foo_id' filter_ = FilterLatestJobIdIsNone() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=None) assert initial_query.count() == 1 assert filtered_query.count() == 0 def test_filter_latest_job_id_is_not_none__apply_included(db, clean_db, apps_unit_of_work): apps_unit_of_work.latest_job_id = None filter_ = FilterLatestJobIdIsNotNone() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=None) assert initial_query.count() == 1 assert filtered_query.count() == 0 def test_filter_latest_job_id_is_not_none__apply_excluded(db, clean_db, apps_unit_of_work): apps_unit_of_work.latest_job_id = 'foo_id' filter_ = FilterLatestJobIdIsNotNone() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=None) assert initial_query.count() == 1 assert filtered_query.count() == 1 def test_filter_report_date_later_than__apply_included(db, clean_db, apps_unit_of_work): report_date = date(2015, 10, 15) filter_value = report_date - timedelta(days=1) apps_unit_of_work.report_date = report_date filter_ = FilterReportDateLaterThan() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=filter_value) assert initial_query.count() == 1 assert filtered_query.count() == 1 def test_filter_report_date_later_than__apply_excluded(db, clean_db, apps_unit_of_work): report_date = date(2015, 10, 15) filter_value = report_date + timedelta(days=1) apps_unit_of_work.report_date = report_date filter_ = FilterReportDateLaterThan() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=filter_value) assert initial_query.count() == 1 assert filtered_query.count() == 0 def test_filter_report_date_earlier_than__apply_included(db, clean_db, apps_unit_of_work): report_date = date(2015, 10, 15) filter_value = report_date + timedelta(days=1) apps_unit_of_work.report_date = report_date filter_ = FilterReportDateEarlierThan() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=filter_value) assert initial_query.count() == 1 assert filtered_query.count() == 1 def test_filter_report_date_earlier_than__apply_excluded(db, clean_db, apps_unit_of_work): report_date = date(2015, 10, 15) filter_value = report_date - timedelta(days=1) apps_unit_of_work.report_date = report_date filter_ = FilterReportDateEarlierThan() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=filter_value) assert initial_query.count() == 1 assert filtered_query.count() == 0 def test_filter_dbx_execution_id_is_none__apply_included( db, clean_db, apps_unit_of_work, dbx_execution ): dbx_execution.status = DBXJobStatusEnum.IN_PROGRESS filter_ = FilterDbxExecutionIdIsNone() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=None) assert initial_query.count() == 1 assert filtered_query.count() == 1 def test_filter_dbx_execution_id_is_none__apply_excluded( db, clean_db, apps_unit_of_work, dbx_execution ): dbx_execution.status = DBXJobStatusEnum.COMPLETE filter_ = FilterDbxExecutionIdIsNone() initial_query = db.query(AppsUnitOfWork) filtered_query = filter_.apply(initial_query, value=None) assert initial_query.count() == 1 assert filtered_query.count() == 0