# pylint: disable=no-member import datetime import unittest.mock import db_schema.factories.apps import db_schema.factories.slz import pytest from dacite.core import from_dict from db_schema import apps from db_schema.schemas.apps import UnitOfWorkPriorityEnum from apps_etl_manager.dependencies.abc import BaseDependency from apps_etl_manager.dependencies.statuses import ( DependentEntities, DependentStatus, DependentStatuses, ) from apps_etl_manager.entities.aws_lambda import ReprocessingPayload from apps_etl_manager.entities.config import UnitOfWorkConfig from apps_etl_manager.entities.entities import EtlArgs, EtlDependency from apps_etl_manager.unit_of_works.uow import EtlUnitOfWork APPLE_TRACKS = UnitOfWorkConfig( unit_of_work='apple-{yyyymmdd}-tracks', schedule='*/10 * * * *', data_source=['slz'], args=EtlArgs( dsp='apple', unit_type='daily', report_type='tracks', valid_from='2020-01-01', dbx_job_id=1, days_offset=1, dbx_job_max_lifetime=1800, ), dependencies=[EtlDependency( type='slz_uow', name='apple-{yyyymmdd}-sme-tracks-v1', )] ) def create_paired_statuses(session): slz_uow = db_schema.factories.slz.UnitOfWorkFactory(unit_of_work_code='unit_of_work_code') sme_streams_uk = DependentStatus( persistent_id=14, licensor_id=1, report_id=3, context='uk', content_name='uk.zip', completed_at=datetime.datetime.now(datetime.timezone.utc), is_single_context=False, slz_uow_id=slz_uow.unit_of_work_id, ) paired_statuses = DependentEntities(session, [sme_streams_uk]) return paired_statuses def get_run_state(uow, session, uow_config=APPLE_TRACKS, reprocessing_payload=None): # build unit and run it unit = EtlUnitOfWork( unittest.mock.Mock(), session, unittest.mock.Mock(), uow, uow_config, 3, datetime.datetime(2020, 8, 12, 1, 2, 3), 'test', reprocessing_payload=reprocessing_payload, ) unit.check_dependencies() run_state = unit.run() return run_state @pytest.mark.integration @unittest.mock.patch( 'apps_etl_manager.unit_of_works.uow.EtlUnitOfWork.get_uow_dependencies', spec=BaseDependency ) def test_etl_uow_check_dependencies(mocked_deps, session): mocked_deps.return_value.is_ready.return_value = True uow = db_schema.factories.apps.UnitOfWorkFactory() unit = EtlUnitOfWork( unittest.mock.Mock(), session, unittest.mock.Mock(), uow, APPLE_TRACKS, 3, datetime.datetime(2020, 8, 12, 1, 2, 3), 'test' ) assert unit.check_dependencies() mocked_deps.return_value.load.assert_called_once() @pytest.mark.integration @unittest.mock.patch( 'apps_etl_manager.unit_of_works.uow.EtlUnitOfWork.get_uow_dependencies', spec=BaseDependency ) def test_etl_uow_run_has_diff_has_unprocessed(mocked_deps, session): """ Tested class: EtlUnitOfWork Diff and unprocceses paired_statuses found, all required deps are completed —> switch to QUEUED status, next_completeness_status is COMPLETE """ # mock unprocessed paired_statuses, next status MIN_COMPLETE mocked_statuses = unittest.mock.MagicMock(spec=DependentStatuses) mocked_statuses.has_unprocessed.return_value = True mocked_statuses.__iter__.return_value = create_paired_statuses(session) # mock new content statuses creation, has_new == True mocked_entities = unittest.mock.MagicMock(spec=DependentEntities) mocked_entities.persist.return_value = True mocked_entities.get_filtered_statuses.return_value = mocked_statuses # mock dependent_entities getting mocked_deps.return_value.dependent_entities.return_value = mocked_entities # mock is_ready for check_dependencies check mocked_deps.return_value.is_ready.return_value = True uow = db_schema.factories.apps.UnitOfWorkFactory() session.commit() # unit will run sf run_state = get_run_state(uow, session) mocked_entities.persist.assert_called_once() assert run_state.sf_name assert run_state.sf_state assert not run_state.is_force_completed assert run_state.sf_state['next_completeness_status' ] == apps.CompletenessStatusEnum.COMPLETE.value assert uow.completeness_status == apps.CompletenessStatusEnum.QUEUED @pytest.mark.integration @unittest.mock.patch( 'apps_etl_manager.unit_of_works.uow.EtlUnitOfWork.get_uow_dependencies', spec=BaseDependency ) def test_etl_uow_run_is_force_completed(mocked_deps, session): """ Tested class: EtlUnitOfWork All uow dependencies are completed, uow have never been failed, there is no diff found for uow —> return is_force_completed=True, will change uow completeness_status to COMPLETE, don’t switch to QUEUED status """ # mock unprocessed paired_statuses, next status COMPLETE mocked_statuses = unittest.mock.MagicMock(spec=DependentStatuses) mocked_statuses.has_unprocessed.return_value = False mocked_statuses.__iter__.return_value = create_paired_statuses(session) # mock new content statuses creation, has_new == False mocked_entities = unittest.mock.MagicMock(spec=DependentEntities) mocked_entities.persist.return_value = False mocked_entities.get_filtered_statuses.return_value = mocked_statuses # mock dependent_entities getting mocked_deps.return_value.dependent_entities.return_value = mocked_entities # mock is_ready for check_dependencies check mocked_deps.return_value.is_ready.return_value = True uow = db_schema.factories.apps.UnitOfWorkFactory() # unit will not run sf, is_force_completed == True run_state = get_run_state(uow, session) assert not run_state.sf_name assert not run_state.sf_state assert run_state.is_force_completed @pytest.mark.integration @unittest.mock.patch( 'apps_etl_manager.unit_of_works.uow.EtlUnitOfWork.get_uow_dependencies', spec=BaseDependency ) def test_etl_uow_run_is_reprocessed(mocked_deps, session): """ Tested class: EtlUnitOfWork All uow dependencies are completed, uow have never been failed, there is no diff found for uow, reprocessing_payload passed —> return is_force_completed=False, next_completeness_status is COMPLETE, UoW switched to QUEUED status """ # mock unprocessed paired_statuses, next status COMPLETE mocked_statuses = unittest.mock.MagicMock(spec=DependentStatuses) mocked_statuses.has_unprocessed.return_value = True mocked_statuses.__iter__.return_value = create_paired_statuses(session) # mock new content statuses creation, has_new == False mocked_entities = unittest.mock.MagicMock(spec=DependentEntities) mocked_entities.persist.return_value = False mocked_entities.get_filtered_statuses.return_value = mocked_statuses # mock dependent_entities getting mocked_deps.return_value.dependent_entities.return_value = mocked_entities # mock is_ready for check_dependencies check mocked_deps.return_value.is_ready.return_value = True uow = db_schema.factories.apps.UnitOfWorkFactory() uow.data_source = apps.DataSourceEnum.SLZ dbx_job_max_lifetime = 444 reprocessing_payload = from_dict( data_class=ReprocessingPayload, data={ 'unit_of_work_id': uow.unit_of_work_id, # pylint: disable=no-member 'priority': UnitOfWorkPriorityEnum.DEFAULT, 'dbx_job_max_lifetime': dbx_job_max_lifetime, } ) run_state = get_run_state(uow, session, reprocessing_payload=reprocessing_payload) assert run_state.sf_name assert run_state.sf_state assert run_state.sf_state['databricks_job_max_lifetime'] == dbx_job_max_lifetime assert not run_state.is_force_completed assert uow.completeness_status == apps.CompletenessStatusEnum.QUEUED @pytest.mark.integration @unittest.mock.patch( 'apps_etl_manager.unit_of_works.uow.EtlUnitOfWork.get_uow_dependencies', spec=BaseDependency ) def test_etl_uow_run_is_unit_failed(mocked_deps, session): """ Tested class: EtlUnitOfWork All uow dependencies are completed, uow have status “FAILED” —> switch to QUEUED status to reprocess uow """ # mock unprocessed paired_statuses, next status COMPLETE mocked_statuses = unittest.mock.MagicMock(spec=DependentStatuses) mocked_statuses.has_unprocessed.return_value = False mocked_statuses.__iter__.return_value = create_paired_statuses(session) # mock new content statuses creation, has_new == False mocked_entities = unittest.mock.MagicMock(spec=DependentEntities) mocked_entities.persist.return_value = False mocked_entities.get_filtered_statuses.return_value = mocked_statuses # mock dependent_entities getting mocked_deps.return_value.dependent_entities.return_value = mocked_entities # mock is_ready for check_dependencies check mocked_deps.return_value.is_ready.return_value = True uow = db_schema.factories.apps.UnitOfWorkFactory() uow.completeness_status = apps.CompletenessStatusEnum.FAILED # unit will not switch to QUEUED status, is_force_completed == True run_state = get_run_state(uow, session) mocked_entities.persist.assert_called_once() assert run_state.sf_name assert run_state.sf_state assert not run_state.is_force_completed assert uow.completeness_status != apps.CompletenessStatusEnum.COMPLETE @pytest.mark.integration @unittest.mock.patch( 'apps_etl_manager.unit_of_works.uow.EtlUnitOfWork.get_uow_dependencies', spec=BaseDependency ) def test_etl_uow_run_has_no_diff_do_nothing(mocked_deps, session): """ Tested class: EtlUnitOfWork Not all uow dependencies are completed, next_status is “MIN_COMPLETE”, uow have never been failed, there is no diff found for uow —> do nothing, don’t switch to QUEUED status """ # mock unprocessed paired_statuses, next status COMPLETE mocked_statuses = unittest.mock.MagicMock(spec=DependentStatuses) mocked_statuses.has_unprocessed.return_value = False mocked_statuses.__iter__.return_value = create_paired_statuses(session) # mock new content statuses creation, has_new == False mocked_entities = unittest.mock.MagicMock(spec=DependentEntities) mocked_entities.persist.return_value = False mocked_entities.get_filtered_statuses.return_value = mocked_statuses # mock dependent_entities getting mocked_deps.return_value.dependent_entities.return_value = mocked_entities # mock is_ready for check_dependencies check mocked_deps.return_value.is_ready.return_value = True # not all uow dependencies are completed mocked_deps.return_value.is_completed.return_value = False uow = db_schema.factories.apps.UnitOfWorkFactory() # unit will not run sf, is_force_completed == True run_state = get_run_state(uow, session) assert not run_state.sf_name assert not run_state.sf_state @pytest.mark.integration @unittest.mock.patch( 'apps_etl_manager.unit_of_works.uow.EtlUnitOfWork.get_uow_dependencies', spec=BaseDependency ) def test_etl_uow_run_with_empty_content_status_ids(mocked_deps, session): """ Tested class: EtlUnitOfWork Uow has_only_single_context_deps next_status is “MIN_COMPLETE”, uow have never been failed —> switch to QUEUED status with empty content_status_ids for has_only_single_context_deps uow """ # mock unprocessed paired_statuses, next status MIN_COMPLETE mocked_statuses = unittest.mock.MagicMock(spec=DependentStatuses) mocked_statuses.has_unprocessed.return_value = True mocked_statuses.__iter__.return_value = create_paired_statuses(session) # mock new content statuses creation, has_new == True mocked_entities = unittest.mock.MagicMock(spec=DependentEntities) mocked_entities.persist.return_value = True mocked_entities.get_filtered_statuses.return_value = mocked_statuses # mock dependent_entities getting mocked_deps.return_value.dependent_entities.return_value = mocked_entities # mock is_ready for check_dependencies check mocked_deps.return_value.is_ready.return_value = True # not all uow dependencies are completed mocked_deps.return_value.is_completed.return_value = False uow = db_schema.factories.apps.UnitOfWorkFactory() uow.data_source = apps.DataSourceEnum.SLZ uow_config = UnitOfWorkConfig( unit_of_work='apple-{yyyymmdd}-tracks', schedule='*/10 * * * *', args=EtlArgs( dsp='apple', unit_type='daily', report_type='tracks', valid_from='2020-01-01', dbx_job_id=1, days_offset=1, dbx_job_max_lifetime=1800, ), dependencies=[ EtlDependency( type='slz_uow', name='apple-{yyyymmdd}-sme-tracks-v1', is_single_context=True ) ], data_source='slz', ) # unit will run sf with empty content_status_ids run_state = get_run_state(uow, session, uow_config) assert run_state.sf_name assert run_state.sf_state assert run_state.sf_state['content_status_ids'] == [] @pytest.mark.integration @unittest.mock.patch( 'apps_etl_manager.unit_of_works.uow.EtlUnitOfWork.get_uow_dependencies', spec=BaseDependency ) def test_etl_uow_run_with_empty_content_status_ids_not_queued(mocked_deps, session): """ Tested class: EtlUnitOfWork Uow not has_only_single_context_deps next_status is “MIN_COMPLETE”, uow have never been failed, no paired_statuses found, no paired_content_status_ids —> do nothing, don’t switch to QUEUED status with empty content_status_ids for not has_only_single_context_deps uow """ # mock unprocessed paired_statuses, next status MIN_COMPLETE mocked_statuses = unittest.mock.MagicMock(spec=DependentStatuses) mocked_statuses.has_unprocessed.return_value = True # no paired_statuses found mocked_statuses.__iter__.return_value = [] # mock new content statuses creation, has_new == True mocked_entities = unittest.mock.MagicMock(spec=DependentEntities) mocked_entities.persist.return_value = True mocked_entities.get_filtered_statuses.return_value = mocked_statuses # mock dependent_entities getting mocked_deps.return_value.dependent_entities.return_value = mocked_entities # mock is_ready for check_dependencies check mocked_deps.return_value.is_ready.return_value = True # not all uow dependencies are completed mocked_deps.return_value.is_completed.return_value = False uow = db_schema.factories.apps.UnitOfWorkFactory() # unit will not run sf with empty content_status_ids run_state = get_run_state(uow, session) assert not run_state.sf_name assert not run_state.sf_state @pytest.mark.integration @unittest.mock.patch( 'apps_etl_manager.unit_of_works.uow.EtlUnitOfWork.get_uow_dependencies', spec=BaseDependency ) def test_etl_uow_run_deps_are_not_completed_set_min_complete(mocked_deps, session): """ Tested class: EtlUnitOfWork Diff and unprocceses paired_statuses found but some of deps are not completed —> switch to QUEUED status, next_completeness_status is MIN_COMPLETE """ # mock unprocessed paired_statuses, next status MIN_COMPLETE mocked_statuses = unittest.mock.MagicMock(spec=DependentStatuses) mocked_statuses.has_unprocessed.return_value = True mocked_statuses.__iter__.return_value = create_paired_statuses(session) # mock new content statuses creation, has_new == True mocked_entities = unittest.mock.MagicMock(spec=DependentEntities) mocked_entities.persist.return_value = True mocked_entities.get_filtered_statuses.return_value = mocked_statuses # mock dependent_entities getting mocked_deps.return_value.dependent_entities.return_value = mocked_entities # mock is_ready for check_dependencies check mocked_deps.return_value.is_ready.return_value = True mocked_deps.return_value.is_completed.return_value = False uow = db_schema.factories.apps.UnitOfWorkFactory() # unit will run sf run_state = get_run_state(uow, session) mocked_entities.persist.assert_called_once() assert run_state.sf_name assert run_state.sf_state assert not run_state.is_force_completed assert run_state.sf_state['next_completeness_status' ] == apps.CompletenessStatusEnum.MIN_COMPLETE.value assert uow.completeness_status == apps.CompletenessStatusEnum.QUEUED @pytest.mark.integration def test_etl_uow_get_conversion_params(monkeypatch, session): monkeypatch.setenv('ENVIRONMENT', 'test') # pylint: disable=line-too-long expected_result = { 'contexts': [{ 'report_type': 'active_claims', 'uow_id': 'youtubereporting-20201206-sme-active_claims-a1-rp20200723T142518', 'compressed_path': 's3://test-sme-data-archive/youtubereporting/active_claims/a1/report_date=2020-12-06/report_licensor=sme/youtubereporting_20201206_content_owner_asset.csv.gz', 'decompressed_bucket': 'test-delphi-sme-data-parquet', 'decompressed_folder': 'youtubereporting/active_claims/a1/report_date=2020-12-06/report_licensor=sme', 'content_name': 'youtubereporting_20201206_content_owner_asset.csv.gz', 'context': 'active_claims' }] } slz_uow = db_schema.factories.slz.UnitOfWorkFactory( unit_of_work_code='youtubereporting-20201206-sme-active_claims-a1', completeness_status='COMPLETE', ) slz_uow.reprocess_id = '20200723T142518' dependent_entities = DependentEntities( session, [ DependentStatus( persistent_id=14, licensor_id=1, report_id=3, context='active_claims', content_name='youtubereporting_20201206_content_owner_asset.csv.gz', completed_at=datetime.datetime.now(datetime.timezone.utc), is_single_context=True, slz_uow_id=slz_uow.unit_of_work_id, ) ] ) etl_uow = db_schema.factories.apps.UnitOfWorkFactory( unit_of_work_code='youtubereporting-20201206-active_claims_conversion', report_date=datetime.date(2020, 12, 6), ) uow_config = UnitOfWorkConfig( unit_of_work='youtubereporting-{yyyymmdd}-active_claims_conversion', schedule='*/10 * * * *', frequency='0 0 * * 7', args=EtlArgs( dsp='youtubereporting', unit_type='weekly', report_type='youtubereporting_active_claims_conversion', valid_from='2020-01-01', dbx_job_id=1, days_offset=0, dbx_job_max_lifetime=1800, ), dependencies=[ EtlDependency( type='slz_uow', name='youtubereporting-{yyyymmdd}-sme-active_claims-a1', is_single_context=True, ) ], data_source='slz', ) unit = EtlUnitOfWork( unittest.mock.Mock(), session, unittest.mock.Mock(), etl_uow, uow_config, 3, datetime.datetime(2020, 12, 23, 1, 2, 3), 'test', ) conversion_params = unit.get_conversion_params(dependent_entities) assert conversion_params == expected_result