import datetime import unittest.mock import pytest from db_schema import apps from db_schema.factories import apps as factories from apps_etl_manager import const from apps_etl_manager.entities.aws_lambda import LambdaConfig, LambdaPayload, ReprocessingPayload from apps_etl_manager.entities.config import UnitOfWorkConfig, UnitOfWorkConfigs from apps_etl_manager.entities.entities import EtlArgs, EtlDependency from apps_etl_manager.unit_of_works.base import BaseUnitOfWork from apps_etl_manager.unit_of_works.uows import UnitOfWorks NOW = datetime.datetime.now(datetime.timezone.utc) MIN_DECREASE_DATETIME = NOW - datetime.timedelta( minutes=24 * 60 * const.PRIORITY_DECREASE_PERIOD_DAYS + 1 ) def create_uow_config(): return UnitOfWorkConfig( unit_of_work="apple-{yyyymmdd}-streams", schedule="*/10 * * * *", data_source="slz", args=EtlArgs( dsp="apple", unit_type="daily", report_type="streams", valid_from="2020-01-01", dbx_job_id=1, days_offset=1, dbx_spark_params=["--driver-java-options", "-Dfoo"], dbx_job_max_lifetime=1800, ), dependencies=[EtlDependency( type="slz_uow", name="apple-{yyyymmdd}-sme-streams-v1", )], ) def create_lambda_config(**extra): kwargs = { "config_path": "config_path", "max_failure_count": 3, "batch_size": 1, "priority_decrease_period_days": 2, "sentry_secret_key": "sentry_secret_key", "rds_secret_key": "rds_secret_key", "sf_secret_key": "sf_secret_key", } if extra: kwargs.update(extra) return LambdaConfig(**kwargs) def create_uow(): return factories.UnitOfWorkFactory( completeness_status="ACTIVE", activity_status="NOT_IN_PROGRESS", unit_of_work_code="apple-20200201-streams", report_date=datetime.date(2020, 2, 1), next_run_at=datetime.datetime(2020, 2, 1, 10), ) @pytest.mark.integration @unittest.mock.patch("apps_etl_manager.unit_of_works.uows.UnitOfWorks.build_unit", autospec=True) def test_run(mocked_build, session): configs = UnitOfWorkConfigs([create_uow_config()]) lambda_config = create_lambda_config() lambda_payload = LambdaPayload( initialize=False, current_dt=datetime.datetime(2020, 2, 2, 10), data_sources=["slz"], ) create_uow() mocked_unit = unittest.mock.Mock(spec=BaseUnitOfWork) mocked_unit.check_dependencies.return_value = True mocked_build.return_value = mocked_unit uows = UnitOfWorks(unittest.mock.Mock(), session, unittest.mock.Mock(), configs, ["slz"]) uows.run(lambda_config, lambda_payload, "test") mocked_unit.check_dependencies.assert_called_once() mocked_unit.run.assert_called_once() @pytest.mark.integration def test_decrease_priorities_update_uows(session): """ Checks that uow's priority will be decreased for UoWs in ACTIVE and MIN_COMPLETE status """ uow1 = create_uow() uow1.created_at = MIN_DECREASE_DATETIME uow2 = create_uow() uow2.completeness_status = apps.CompletenessStatusEnum.MIN_COMPLETE uow2.created_at = MIN_DECREASE_DATETIME configs = UnitOfWorkConfigs([create_uow_config()]) uows = UnitOfWorks(unittest.mock.Mock(), session, unittest.mock.Mock(), configs, ["slz"]) lambda_config = create_lambda_config() uows.decrease_priorities(lambda_config) assert uow1.priority == 8 assert uow2.priority == 8 @pytest.mark.integration @pytest.mark.parametrize( "created_at, priority", ( (NOW, 8), (NOW, 5), (MIN_DECREASE_DATETIME, 8), (MIN_DECREASE_DATETIME, 4), ), ) def test_decrease_priorities_do_not_update_uows(session, created_at, priority): """ Checks that uow's priority will not be decreased """ uow = create_uow() uow.created_at = created_at uow.priority = priority configs = UnitOfWorkConfigs([create_uow_config()]) uows = UnitOfWorks(unittest.mock.Mock(), session, unittest.mock.Mock(), configs, ["slz"]) lambda_config = create_lambda_config() uows.decrease_priorities(lambda_config) assert uow.priority == priority @pytest.mark.integration @pytest.mark.parametrize( "status", [apps.CompletenessStatusEnum.FAILED, apps.CompletenessStatusEnum.COMPLETE] ) def test_reprocess_run_started(mocker, session, status): configs = UnitOfWorkConfigs([create_uow_config()]) unit_of_work: apps.UnitOfWork = create_uow() unit_of_work.completeness_status = status new_priority = apps.UnitOfWorkPriorityEnum.PRIORITY_4 lambda_config = create_lambda_config() lambda_payload = LambdaPayload( initialize=False, current_dt=datetime.datetime(2020, 2, 2, 10), data_sources=["slz"], reprocessing=[ ReprocessingPayload( unit_of_work_id=unit_of_work.unit_of_work_id, # pylint: disable=no-member priority=new_priority, ) ], ) mocked_unit = mocker.Mock(spec=BaseUnitOfWork) mocked_unit.check_dependencies.return_value = True mocked_build_unit = mocker.patch( "apps_etl_manager.unit_of_works.uows.UnitOfWorks.build_unit", return_value=mocked_unit, autospec=True, ) uows = UnitOfWorks(unittest.mock.Mock(), session, unittest.mock.Mock(), configs, ["slz"]) uows.run(lambda_config, lambda_payload, "test") mocked_unit.check_dependencies.assert_called_once() mocked_unit.run.assert_called_once() assert mocked_build_unit.call_args[1]["reprocessing_payload"].priority == new_priority @pytest.mark.integration def test_reprocess_run_multiple_started(mocker, session): configs = UnitOfWorkConfigs([create_uow_config()]) count = 5 unit_of_works = [] for _ in range(0, count): unit_of_work: apps.UnitOfWork = create_uow() unit_of_work.completeness_status = apps.CompletenessStatusEnum.FAILED unit_of_work.activity_status = apps.ActivityStatusEnum.NOT_IN_PROGRESS unit_of_works.append(unit_of_work) lambda_config = create_lambda_config(batch_size=count) lambda_payload = LambdaPayload( initialize=False, current_dt=datetime.datetime(2020, 2, 2, 10), data_sources=["slz"], reprocessing=[ ReprocessingPayload( unit_of_work_id=unit_of_work.unit_of_work_id, # pylint: disable=no-member priority=apps.UnitOfWorkPriorityEnum.PRIORITY_4, ) for unit_of_work in unit_of_works ], ) mocked_unit = mocker.Mock(spec=BaseUnitOfWork) mocked_unit.check_dependencies.return_value = True mocker.patch( "apps_etl_manager.unit_of_works.uows.UnitOfWorks.build_unit", return_value=mocked_unit, autospec=True, ) uows = UnitOfWorks(unittest.mock.Mock(), session, unittest.mock.Mock(), configs, ["slz"]) uows.run(lambda_config, lambda_payload, "test") assert mocked_unit.check_dependencies.call_count == count assert mocked_unit.run.call_count == count @pytest.mark.integration def test_reprocess_remain_the_same_priority(mocker, session): configs = UnitOfWorkConfigs([create_uow_config()]) unit_of_work: apps.UnitOfWork = create_uow() unit_of_work.completeness_status = apps.CompletenessStatusEnum.FAILED lambda_config = create_lambda_config() lambda_payload = LambdaPayload( initialize=False, current_dt=datetime.datetime(2020, 2, 2, 10), data_sources=["slz"], reprocessing=[ReprocessingPayload(unit_of_work_id=unit_of_work.unit_of_work_id, )], # pylint: disable=no-member ) mocked_unit = mocker.Mock(spec=BaseUnitOfWork) mocked_unit.check_dependencies.return_value = True mocked_build_unit = mocker.patch( "apps_etl_manager.unit_of_works.uows.UnitOfWorks.build_unit", return_value=mocked_unit, autospec=True, ) uows = UnitOfWorks(unittest.mock.Mock(), session, unittest.mock.Mock(), configs, ["slz"]) uows.run(lambda_config, lambda_payload, "test") mocked_unit.check_dependencies.assert_called_once() mocked_unit.run.assert_called_once() assert mocked_build_unit.call_args[1]["reprocessing_payload"].priority is None @pytest.mark.integration @pytest.mark.parametrize( "status", [ x for x in apps.CompletenessStatusEnum if x not in ( apps.CompletenessStatusEnum.FAILED, apps.CompletenessStatusEnum.COMPLETE, ) ], ) def test_reprocess_run_skipped_for_invalid_completeness_status(mocker, session, status): configs = UnitOfWorkConfigs([create_uow_config()]) unit_of_work: apps.UnitOfWork = create_uow() unit_of_work.completeness_status = status unit_of_work.activity_status = apps.ActivityStatusEnum.NOT_IN_PROGRESS lambda_config = create_lambda_config() lambda_payload = LambdaPayload( initialize=False, current_dt=datetime.datetime(2020, 2, 2, 10), data_sources=["slz"], reprocessing=[ ReprocessingPayload( unit_of_work_id=unit_of_work.unit_of_work_id, # pylint: disable=no-member priority=apps.UnitOfWorkPriorityEnum.PRIORITY_4, ) ], ) mocked_unit = mocker.Mock(spec=BaseUnitOfWork) mocked_unit.check_dependencies.return_value = True mocker.patch( "apps_etl_manager.unit_of_works.uows.UnitOfWorks.build_unit", return_value=mocked_unit, autospec=True, ) uows = UnitOfWorks(unittest.mock.Mock(), session, unittest.mock.Mock(), configs, ["slz"]) uows.run(lambda_config, lambda_payload, "test") assert not mocked_unit.check_dependencies.called assert not mocked_unit.run.called @pytest.mark.integration @pytest.mark.parametrize( "status", [x for x in apps.ActivityStatusEnum if x != apps.ActivityStatusEnum.NOT_IN_PROGRESS], ) def test_reprocess_run_skipped_for_invalid_activity_status(mocker, session, status): configs = UnitOfWorkConfigs([create_uow_config()]) unit_of_work: apps.UnitOfWork = create_uow() unit_of_work.completeness_status = apps.CompletenessStatusEnum.FAILED unit_of_work.activity_status = status lambda_config = create_lambda_config() lambda_payload = LambdaPayload( initialize=False, current_dt=datetime.datetime(2020, 2, 2, 10), data_sources=["slz"], reprocessing=[ ReprocessingPayload( unit_of_work_id=unit_of_work.unit_of_work_id, # pylint: disable=no-member priority=apps.UnitOfWorkPriorityEnum.PRIORITY_4, ) ], ) mocked_unit = mocker.Mock(spec=BaseUnitOfWork) mocked_unit.check_dependencies.return_value = True mocker.patch( "apps_etl_manager.unit_of_works.uows.UnitOfWorks.build_unit", return_value=mocked_unit, autospec=True, ) uows = UnitOfWorks(unittest.mock.Mock(), session, unittest.mock.Mock(), configs, ["slz"]) uows.run(lambda_config, lambda_payload, "test") assert not mocked_unit.check_dependencies.called assert not mocked_unit.run.called @pytest.mark.integration def test_reprocess_run_skipped_for_data_sources_mismatch(mocker, session): configs = UnitOfWorkConfigs([create_uow_config()]) unit_of_work: apps.UnitOfWork = create_uow() unit_of_work.completeness_status = apps.CompletenessStatusEnum.FAILED unit_of_work.activity_status = apps.ActivityStatusEnum.NOT_IN_PROGRESS lambda_config = create_lambda_config() lambda_payload = LambdaPayload( initialize=False, current_dt=datetime.datetime(2020, 2, 2, 10), data_sources=["slz"], reprocessing=[ ReprocessingPayload( unit_of_work_id=unit_of_work.unit_of_work_id, # pylint: disable=no-member priority=apps.UnitOfWorkPriorityEnum.PRIORITY_4, ) ], ) mocked_unit = mocker.Mock(spec=BaseUnitOfWork) mocked_unit.check_dependencies.return_value = True mocker.patch( "apps_etl_manager.unit_of_works.uows.UnitOfWorks.build_unit", return_value=mocked_unit, autospec=True, ) uows = UnitOfWorks( unittest.mock.Mock(), session, unittest.mock.Mock(), configs, ["chartmetric"] ) uows.run(lambda_config, lambda_payload, "test") assert not mocked_unit.check_dependencies.called assert not mocked_unit.run.called @pytest.mark.integration @pytest.mark.parametrize( ("is_reprocessing", "add_params", "override_params", "result"), ( (True, ["foo", "bar"], None, ["--driver-java-options", "-Dfoo foo bar"]), (True, None, ["foo", "bar"], ["--driver-java-options", "foo bar"]), (True, ["baz", "foobar"], ["foo", "bar"], ["--driver-java-options", "foo bar"]), (False, ["baz", "foobar"], ["foo", "bar"], ["--driver-java-options", "-Dfoo"]), ), ) # yapf: disable def test_reprocess_dbx_params( mocker, session, is_reprocessing, add_params, override_params, result ): configs = UnitOfWorkConfigs([create_uow_config()]) unit_of_work: apps.UnitOfWork = create_uow() unit_of_work.completeness_status = apps.CompletenessStatusEnum.FAILED unit_of_work.activity_status = apps.ActivityStatusEnum.NOT_IN_PROGRESS lambda_config = create_lambda_config() lambda_payload = LambdaPayload( initialize=False, current_dt=datetime.datetime(2020, 2, 2, 10), data_sources=["slz"], reprocessing=[ ReprocessingPayload( unit_of_work_id=unit_of_work.unit_of_work_id, # pylint: disable=no-member priority=apps.UnitOfWorkPriorityEnum.PRIORITY_4, ) ] if is_reprocessing else None, add_driver_java_options=add_params, override_driver_java_options=override_params, ) mocked_build_unit = mocker.patch( "apps_etl_manager.unit_of_works.uows.UnitOfWorks.build_unit", autospec=True, ) uows = UnitOfWorks(unittest.mock.Mock(), session, unittest.mock.Mock(), configs, ["slz"]) uows.run(lambda_config, lambda_payload, "test") assert mocked_build_unit.called config = mocked_build_unit.call_args[1]["uow_cfg"] assert config.args.dbx_spark_params == result @pytest.mark.integration @pytest.mark.parametrize( ("failure_count", "completeness_status", "is_uow_in_query"), ( (1, apps.CompletenessStatusEnum.FAILED, True), (2, apps.CompletenessStatusEnum.FAILED, True), (3, apps.CompletenessStatusEnum.MIN_COMPLETE, True), (4, apps.CompletenessStatusEnum.ACTIVE, True), (3, apps.CompletenessStatusEnum.FAILED, False), (0, apps.CompletenessStatusEnum.COMPLETE, False), (0, apps.CompletenessStatusEnum.QUEUED, False), ), ) # yapf: disable def test_build_query_contains_uow(session, failure_count, completeness_status, is_uow_in_query): uow = create_uow() uow.completeness_status = completeness_status uow.failure_count = failure_count configs = UnitOfWorkConfigs([create_uow_config()]) uows = UnitOfWorks(unittest.mock.Mock(), session, unittest.mock.Mock(), configs, ["slz"]) lambda_config = create_lambda_config() lambda_payload = LambdaPayload( initialize=False, current_dt=datetime.datetime(2020, 2, 2, 10), data_sources=["slz"], reprocessing=None, ) query = uows.build_query(lambda_config, lambda_payload) assert (uow in query) == is_uow_in_query @pytest.mark.integration @pytest.mark.parametrize( ("report_date", "dependency_name"), ( (datetime.date(2021, 4, 20), "apple-{yyyymmdd}-sme-amSummaryStreams-v1_0"), (datetime.date(2020, 3, 14), "apple-{yyyymmdd}-sme-amStreams-v1_2"), ), ) def test_build_config_double_uows(session, report_date, dependency_name): unit_of_work_code = "apple-{yyyymmdd}-streams_isrc_date_day" uow_configs = UnitOfWorkConfigs([ UnitOfWorkConfig( unit_of_work=unit_of_work_code, schedule="*/10 * * * *", data_source="slz", args=EtlArgs( dsp="apple", unit_type="daily", report_type="apple_streams_isrc_date_day", valid_from="2021-04-14", dbx_job_id=1, days_offset=1, dbx_spark_params=["--driver-java-options", "-Dfoo"], dbx_job_max_lifetime=1800, ), dependencies=[ EtlDependency( type="slz_uow", name="apple-{yyyymmdd}-sme-amSummaryStreams-v1_0", ) ], ), UnitOfWorkConfig( unit_of_work=unit_of_work_code, schedule="*/10 * * * *", data_source="slz", args=EtlArgs( dsp="apple", unit_type="daily", report_type="apple_streams_isrc_date_day", valid_from="2017-06-30", valid_until="2021-04-13", dbx_job_id=1, days_offset=1, dbx_spark_params=["--driver-java-options", "-Dfoo"], dbx_job_max_lifetime=1800, ), dependencies=[ EtlDependency( type="slz_uow", name="apple-{yyyymmdd}-sme-amStreams-v1_2", ) ], ), ]) lambda_payload = LambdaPayload( initialize=True, current_dt=datetime.datetime.now(), data_sources=["slz"] ) uows = UnitOfWorks(unittest.mock.Mock(), session, unittest.mock.Mock(), uow_configs, ["slz"]) uow = factories.UnitOfWorkFactory( completeness_status="ACTIVE", activity_status="NOT_IN_PROGRESS", unit_of_work_code=unit_of_work_code, report_date=report_date, ) uow_config = uows.build_config(lambda_payload, uow) assert uow_config.dependencies[0].name == dependency_name