import datetime import unittest.mock import pytest from db_schema import apps, common from db_schema.factories import apps as factories from apps_etl_runner.entities import LambdaConfig from apps_etl_runner.pool import SimpleRunner, UowPool def create_unit(priority=5, data_source='SLZ'): return factories.UnitOfWorkFactory( activity_status=common.ActivityStatusEnum.NOT_IN_PROGRESS, completeness_status='QUEUED', latest_job_id='sf-name', latest_job_state={'dbx_job_id': 14}, priority=priority, data_source=data_source ) def create_lambda_config( max_active_runs=2, backfill_active_runs=0, priority_threshold=8, process_backfill_uows=False ): return LambdaConfig( etl_sf_arn='test', chartmetric_sf_arn='test', max_active_runs=max_active_runs, backfill_active_runs=backfill_active_runs, priority_threshold=priority_threshold, process_backfill_uows=process_backfill_uows, dbx_job_max_lifetime=1800, sentry_secret_key='sentry_secret_key', rds_secret_key='rds_secret_key' ) @pytest.mark.integration def test_uow_pool_load(session): factories.UnitOfWorkFactory(activity_status=common.ActivityStatusEnum.IN_PROGRESS) factories.UnitOfWorkFactory(activity_status=common.ActivityStatusEnum.NOT_IN_PROGRESS) factories.UnitOfWorkFactory( activity_status=common.ActivityStatusEnum.IN_PROGRESS, data_source='CHARTMETRIC', ) factories.UnitOfWorkFactory( activity_status=common.ActivityStatusEnum.IN_PROGRESS, unit_of_work_type=apps.UnitOfWorkTypeEnum.BACKFILL ) config = create_lambda_config() pool = UowPool(unittest.mock.Mock(), session, unittest.mock.Mock(), config) pool.load() assert pool._num_of_runs == 1 @pytest.mark.integration def test_uow_pool_can_run(session): factories.UnitOfWorkFactory(activity_status=common.ActivityStatusEnum.IN_PROGRESS) factories.UnitOfWorkFactory(activity_status=common.ActivityStatusEnum.IN_PROGRESS) config = create_lambda_config(max_active_runs=3) pool = UowPool(unittest.mock.Mock(), session, unittest.mock.Mock(), config) pool.load() assert pool.can_run() factories.UnitOfWorkFactory(activity_status=common.ActivityStatusEnum.IN_PROGRESS) pool.load() assert not pool.can_run() @pytest.mark.integration def test_uow_pool_increment(session): config = create_lambda_config() pool = UowPool(unittest.mock.Mock(), session, unittest.mock.Mock(), config) pool.load() assert pool.can_run() pool.increment() assert pool.can_run() pool.increment() assert not pool.can_run() @pytest.mark.integration def test_uow_pool_run_max_active_runs(session): config = create_lambda_config(max_active_runs=1) pool = UowPool(unittest.mock.Mock(), session, unittest.mock.Mock(), config) factories.UnitOfWorkFactory(activity_status=common.ActivityStatusEnum.IN_PROGRESS) factories.UnitOfWorkFactory(activity_status=common.ActivityStatusEnum.IN_PROGRESS) unit = create_unit(priority=1) sf_names = pool.run() assert unit.activity_status == common.ActivityStatusEnum.NOT_IN_PROGRESS assert len(sf_names) == 0 @pytest.mark.integration def test_uow_pool_priority_threshold(session): """ Test query filtering: exlude backfill uows """ config = create_lambda_config(max_active_runs=2) pool = UowPool(unittest.mock.Mock(), session, unittest.mock.Mock(), config) unit1 = create_unit(priority=9) unit2 = create_unit(priority=8) unit3 = create_unit(priority=5) sf_names = pool.run() assert unit1.activity_status == common.ActivityStatusEnum.NOT_IN_PROGRESS assert unit2.activity_status == common.ActivityStatusEnum.IN_PROGRESS assert unit3.activity_status == common.ActivityStatusEnum.IN_PROGRESS assert len(sf_names) == 2 @pytest.mark.integration def test_uow_pool_disable_process_backfill_uows(session): """ Test query filtering: exlude backfill uows """ config = create_lambda_config(max_active_runs=1, backfill_active_runs=1, priority_threshold=9) pool = UowPool(unittest.mock.Mock(), session, unittest.mock.Mock(), config) unit1 = create_unit(priority=5) unit2 = create_unit(priority=9) unit2.unit_of_work_type = 'BACKFILL' sf_names = pool.run() assert len(sf_names) == 1 assert unit1.activity_status == common.ActivityStatusEnum.IN_PROGRESS assert unit2.activity_status == common.ActivityStatusEnum.NOT_IN_PROGRESS pool = UowPool(unittest.mock.Mock(), session, unittest.mock.Mock(), config, True) sf_names = pool.run() assert len(sf_names) == 1 assert unit2.activity_status == common.ActivityStatusEnum.IN_PROGRESS @pytest.mark.integration def test_simple_runner_get_query(session): """ Test query filters """ unit1 = create_unit(priority=1, data_source='SLZ') unit2 = create_unit(priority=1, data_source='CHARTMETRIC') config = create_lambda_config() pool = SimpleRunner(unittest.mock.Mock(), session, unittest.mock.Mock(), config, 'chartmetric') query = pool.get_query() assert unit1 not in query assert unit2 in query @pytest.mark.integration def test_simple_runner_run(session): """ Test sf running """ unit1 = create_unit(priority=1, data_source='SLZ') unit2 = create_unit(priority=1, data_source='CHARTMETRIC') config = create_lambda_config() pool = SimpleRunner(unittest.mock.Mock(), session, unittest.mock.Mock(), config, 'chartmetric') sf_names = pool.run() assert len(sf_names) == 1 assert unit1.activity_status == common.ActivityStatusEnum.NOT_IN_PROGRESS assert unit2.activity_status == common.ActivityStatusEnum.IN_PROGRESS @pytest.mark.integration def test_uow_pool_clean(session): """ Test that pool cleans stucked in IN_PROGRESS status dbx executions and it's UoWs """ current_dt = datetime.datetime.now(datetime.timezone.utc) datetime_in_future = current_dt + datetime.timedelta(seconds=10) datetime_in_past = current_dt - datetime.timedelta(seconds=10) unit1 = factories.UnitOfWorkFactory(activity_status=common.ActivityStatusEnum.IN_PROGRESS) unit2 = factories.UnitOfWorkFactory( activity_status=common.ActivityStatusEnum.IN_PROGRESS, completeness_status=apps.CompletenessStatusEnum.QUEUED, last_updated_at=datetime_in_past, ) execution1 = factories.DatabricksExecutionFactory( unit_of_work=unit1, dbx_job_id=14, status='IN_PROGRESS', dbx_job_expired_at=datetime_in_future, ) execution2 = factories.DatabricksExecutionFactory( unit_of_work=unit2, dbx_job_id=14, status='IN_PROGRESS', dbx_job_expired_at=datetime_in_past, ) session.commit() config = create_lambda_config() pool = UowPool(unittest.mock.Mock(), session, unittest.mock.Mock(), config) pool.load() assert pool._num_of_runs == 2 pool._clean() pool.load() assert pool._num_of_runs == 1 assert execution1.status == apps.DBXJobStatusEnum.IN_PROGRESS assert execution2.status == apps.DBXJobStatusEnum.FAILED assert execution2.completed_at # pylint: disable=E1101 session.refresh(unit2) assert unit1.activity_status == common.ActivityStatusEnum.IN_PROGRESS assert unit2.activity_status == common.ActivityStatusEnum.NOT_IN_PROGRESS assert unit2.completeness_status == apps.CompletenessStatusEnum.FAILED assert unit2.last_updated_at > datetime_in_past @pytest.mark.integration def test_uow_pool_cancel_outdated_snapshot_uows(session): """ Test that pool switchs outdated UoWs to CANCELLED completeness status """ report = factories.ReportFactory() unit1 = factories.UnitOfWorkFactory( activity_status=common.ActivityStatusEnum.NOT_IN_PROGRESS, completeness_status=apps.CompletenessStatusEnum.QUEUED, unit_of_work_type=apps.UnitOfWorkTypeEnum.SNAPSHOT, report_date=datetime.date(2021, 3, 4), report=report ) unit2 = factories.UnitOfWorkFactory( activity_status=common.ActivityStatusEnum.IN_PROGRESS, completeness_status=apps.CompletenessStatusEnum.QUEUED, unit_of_work_type=apps.UnitOfWorkTypeEnum.SNAPSHOT, report_date=datetime.date(2021, 3, 7), report=report ) unit3 = factories.UnitOfWorkFactory( activity_status=common.ActivityStatusEnum.NOT_IN_PROGRESS, completeness_status=apps.CompletenessStatusEnum.ACTIVE, unit_of_work_type=apps.UnitOfWorkTypeEnum.SNAPSHOT, report_date=datetime.date(2021, 3, 10), report=report ) factories.UnitOfWorkFactory( activity_status=common.ActivityStatusEnum.NOT_IN_PROGRESS, completeness_status=apps.CompletenessStatusEnum.ACTIVE, data_source='SLZ', unit_of_work_type=apps.UnitOfWorkTypeEnum.SNAPSHOT, report_date=datetime.date(2021, 3, 14), report=report ) config = create_lambda_config() pool = UowPool(unittest.mock.Mock(), session, unittest.mock.Mock(), config) # newer uow is IN_PROGRESS, outdated uow should become CANCELLED newer_uow_exists = pool._check_snapshot_uow_is_outdated(unit1) assert unit1.activity_status == common.ActivityStatusEnum.NOT_IN_PROGRESS assert unit1.completeness_status == apps.CompletenessStatusEnum.CANCELLED assert newer_uow_exists is True # newer uow in COMPLETE, outdated uow should become CANCELLED unit3.completeness_status = apps.CompletenessStatusEnum.COMPLETE newer_uow_exists = pool._check_snapshot_uow_is_outdated(unit2) assert unit2.activity_status == common.ActivityStatusEnum.NOT_IN_PROGRESS assert unit2.completeness_status == apps.CompletenessStatusEnum.CANCELLED assert newer_uow_exists is True # newer uow isn't IN_PROGRESS or COMPLETE, outdated uow should be launched unit3.completeness_status = apps.CompletenessStatusEnum.QUEUED newer_uow_exists = pool._check_snapshot_uow_is_outdated(unit3) assert unit3.completeness_status != apps.CompletenessStatusEnum.CANCELLED assert newer_uow_exists is False