# pylint: disable=no-member import datetime import unittest.mock import db_schema import pytest from db_schema.factories import apps, slz from apps_etl_manager.const import SLZ_GRAS_UOWS_DELAY from apps_etl_manager.dependencies.dependencies import EtlUoWDependencies, UoWOptimizedDependencies from apps_etl_manager.dependencies.statuses import DependentEntities, DependentStatus, DependentUoW from apps_etl_manager.entities.config import Dependencies from apps_etl_manager.entities.entities import EtlDependency def create_dependent_status( licensor_id: int, report_id: int, context: str, single: bool = False, ): slz_uow = slz.UnitOfWorkFactory(unit_of_work_code='unit_of_work_code') return DependentStatus( licensor_id=licensor_id, report_id=report_id, context=context, content_name=f'{context}.zip', is_single_context=single, completed_at=datetime.datetime.now(datetime.timezone.utc), slz_uow_id=slz_uow.unit_of_work_id, ) def create_depdendent_uow(licensor_id: int, report_id: int, single=False): return DependentUoW( licensor_id=licensor_id, report_id=report_id, is_single_context=single, ) @pytest.mark.integration def test_depdendencies_load_ok(session): apps_uow = apps.UnitOfWorkFactory(report_date=datetime.date(2020, 1, 1), ) slz.UnitOfWorkFactory( unit_of_work_code='apple-20200101-sme-tracks-v1', report_date=datetime.date(2020, 1, 1), completeness_status='COMPLETE', ) uow_dependencies = EtlUoWDependencies( unittest.mock.Mock(), session, unittest.mock.Mock(), dependencies=Dependencies([ EtlDependency(type='slz_uow', name='apple-{yyyymmdd}-sme-tracks-v1'), ]), date=apps_uow.report_date, include_on_hold=False, dsp='apple', env='test', ) uow_dependencies.load() assert uow_dependencies.is_ready() assert uow_dependencies.is_completed() assert len(uow_dependencies.dependent_entities()) == 1 @pytest.mark.integration def test_depdendencies_load_not_ready(session): apps_uow = apps.UnitOfWorkFactory(report_date=datetime.date(2020, 1, 1), ) uow_dependencies = EtlUoWDependencies( unittest.mock.Mock(), session, unittest.mock.Mock(), dependencies=Dependencies([ EtlDependency(type='slz_uow', name='apple-{yyyymmdd}-sme-tracks-v1'), ]), date=apps_uow.report_date, include_on_hold=False, dsp='apple', env='test', ) uow_dependencies.load() assert not uow_dependencies.is_ready() assert not uow_dependencies.is_completed() assert not uow_dependencies.dependent_entities() @pytest.mark.integration def test_depdendencies_load_not_ready_and_not_completed(session): apps_uow = apps.UnitOfWorkFactory(report_date=datetime.date(2020, 1, 1), ) slz.UnitOfWorkFactory( unit_of_work_code='apple-20200101-sme-tracks-v1', report_date=datetime.date(2020, 1, 1), completeness_status='ACTIVE', ) uow_dependencies = EtlUoWDependencies( unittest.mock.Mock(), session, unittest.mock.Mock(), dependencies=Dependencies([ EtlDependency(type='slz_uow', name='apple-{yyyymmdd}-sme-tracks-v1'), ]), date=apps_uow.report_date, include_on_hold=False, dsp='apple', env='test', ) uow_dependencies.load() assert not uow_dependencies.is_ready() assert not uow_dependencies.is_completed() assert len(uow_dependencies.dependent_entities()) == 0 @pytest.mark.integration def test_depdendencies_load_ready_and_not_completed(session): apps_uow = apps.UnitOfWorkFactory(report_date=datetime.date(2020, 1, 1), ) slz_uow = slz.UnitOfWorkFactory( unit_of_work_code='apple-20200101-sme-tracks-v1', report_date=datetime.date(2020, 1, 1), completeness_status='MIN_COMPLETE', ) slz.ContentStatusFactory(unit_of_work=slz_uow, content_status='COMPLETE') slz.ContentStatusFactory(unit_of_work=slz_uow, content_status='ACTIVE') session.commit() uow_dependencies = EtlUoWDependencies( unittest.mock.Mock(), session, unittest.mock.Mock(), dependencies=Dependencies([ EtlDependency(type='slz_uow', name='apple-{yyyymmdd}-sme-tracks-v1'), ]), date=apps_uow.report_date, include_on_hold=False, dsp='apple', env='test', ) uow_dependencies.load() assert uow_dependencies.is_ready() assert not uow_dependencies.is_completed() assert len(uow_dependencies.dependent_entities()) == 1 @pytest.mark.integration def test_depdendencies_include_on_hold_load_ready_and_not_completed(session): apps_uow = apps.UnitOfWorkFactory(report_date=datetime.date(2020, 1, 1), ) slz_uow = slz.UnitOfWorkFactory( unit_of_work_code='apple-20200101-sme-tracks-v1', report_date=datetime.date(2020, 1, 1), completeness_status='MIN_COMPLETE', ) slz.ContentStatusFactory(unit_of_work=slz_uow, content_status='COMPLETE') slz.ContentStatusFactory(unit_of_work=slz_uow, content_status='ACTIVE') slz.ContentStatusFactory(unit_of_work=slz_uow, content_status='ON_HOLD') slz.ContentStatusFactory(unit_of_work=slz_uow, content_status='ON_HOLD', completed_at=None) uow_dependencies = EtlUoWDependencies( unittest.mock.Mock(), session, unittest.mock.Mock(), dependencies=Dependencies([ EtlDependency(type='slz_uow', name='apple-{yyyymmdd}-sme-tracks-v1'), ]), date=apps_uow.report_date, include_on_hold=True, dsp='apple', env='test' ) session.commit() uow_dependencies.load() assert uow_dependencies.is_ready() assert not uow_dependencies.is_completed() assert len(uow_dependencies.dependent_entities()) == 3 @pytest.mark.integration def test_depdendencies_load_if_has_optional(session): apps_uow = apps.UnitOfWorkFactory(report_date=datetime.date(2020, 1, 1), ) slz.UnitOfWorkFactory( unit_of_work_code='apple-20200101-sme-streams-v1', report_date=datetime.date(2020, 1, 1), completeness_status='ACTIVE', ) completed = slz.UnitOfWorkFactory( unit_of_work_code='apple-20200101-sme-tracks-v1', report_date=datetime.date(2020, 1, 1), completeness_status='COMPLETE', ) slz.ContentStatusFactory( unit_of_work=completed, content_status='COMPLETE', ) uow_dependencies = EtlUoWDependencies( unittest.mock.Mock(), session, unittest.mock.Mock(), dependencies=Dependencies([ EtlDependency( type='slz_uow', name='apple-{yyyymmdd}-sme-tracks-v1', ), EtlDependency( type='slz_uow', name='apple-{yyyymmdd}-sme-streams-v1', is_optional=True, ), ]), date=apps_uow.report_date, include_on_hold=False, dsp='apple', env='test', ) session.commit() uow_dependencies.load() assert uow_dependencies.is_ready() assert not uow_dependencies.is_completed() dependent_entities = uow_dependencies.dependent_entities() assert len(dependent_entities) == 2 assert isinstance(dependent_entities[0], DependentStatus) assert isinstance(dependent_entities[1], DependentUoW) @pytest.mark.integration def test_dependent_statuses_for_optional_active_unit(session): apps_uow = apps.UnitOfWorkFactory(report_date=datetime.date(2020, 2, 2), ) uow = slz.UnitOfWorkFactory( unit_of_work_code='apple-20200202-sme-users-v1', completeness_status='ACTIVE', ) slz.ContentStatusFactory( unit_of_work=uow, content_status='COMPLETE', ) slz.ContentStatusFactory( unit_of_work=uow, content_status='ACTIVE', ) uow_dependencies = EtlUoWDependencies( unittest.mock.Mock(), session, unittest.mock.Mock(), dependencies=Dependencies([ EtlDependency( type='slz_uow', name='apple-{yyyymmdd}-sme-users-v1', is_optional=True, ), ]), date=apps_uow.report_date, include_on_hold=False, dsp='apple', env='test', ) uow_dependencies.load() dependent_entities = uow_dependencies.dependent_entities() assert len(dependent_entities) == 1 assert isinstance(dependent_entities[0], DependentUoW) @pytest.mark.integration def test_get_filtered_statuses_only_depdendent_statuses(session): sme_users_single = create_dependent_status( licensor_id=1, report_id=1, context='users', single=True, ) sme_tracks_us = create_dependent_status( licensor_id=1, report_id=2, context='us', ) sme_tracks_uk = create_dependent_status( licensor_id=1, report_id=2, context='uk', ) sme_streams_uk = create_dependent_status( licensor_id=1, report_id=3, context='uk', ) entities = DependentEntities( session, [ sme_users_single, sme_tracks_us, sme_tracks_uk, sme_streams_uk, ] ) assert entities.get_filtered_statuses() == [sme_tracks_uk, sme_streams_uk] @pytest.mark.integration def test_get_filtered_statuses_with_depdendent_uow_not_paired(session): sme_tracks_uk = create_dependent_status( licensor_id=1, report_id=2, context='uk', ) sme_streams_uk = create_dependent_status( licensor_id=1, report_id=3, context='uk', ) theorchard_streams = create_depdendent_uow( licensor_id=2, report_id=3, ) entities = DependentEntities(session, [ theorchard_streams, sme_tracks_uk, sme_streams_uk, ]) assert entities.get_filtered_statuses() == [sme_tracks_uk, sme_streams_uk] @pytest.mark.integration def test_get_filtered_statuses_with_depdendent_uow_paired(session): sme_tracks_uk = create_dependent_status( licensor_id=1, report_id=2, context='uk', ) sme_streams_uk = create_dependent_status( licensor_id=1, report_id=3, context='uk', ) theorchard_tracks_uk = create_dependent_status(licensor_id=2, report_id=2, context='uk') theorchard_streams = create_depdendent_uow( licensor_id=2, report_id=3, ) entities = DependentEntities( session, [ theorchard_tracks_uk, theorchard_streams, sme_tracks_uk, sme_streams_uk, ] ) assert entities.get_filtered_statuses() == [sme_tracks_uk, sme_streams_uk] @pytest.mark.integration def test_persistent_statuses(session): sme = slz.LicensorFactory() tracks = slz.ReportFactory() streams = slz.ReportFactory() uow = apps.UnitOfWorkFactory() sme_tracks_uk = create_dependent_status( licensor_id=sme.licensor_id, report_id=tracks.report_id, context='uk', ) sme_streams_uk = create_dependent_status( licensor_id=sme.licensor_id, report_id=streams.report_id, context='uk', ) depdendent_entities = DependentEntities(session, [ sme_tracks_uk, sme_streams_uk, ]) has_new = depdendent_entities.persist(uow) persistent_ids = {entity.persistent_id for entity in depdendent_entities} content_status_ids = {status.content_status_id for status in uow.content_statuses} assert has_new assert len(uow.content_statuses) == 2 assert persistent_ids == content_status_ids @pytest.mark.integration @pytest.mark.parametrize( 'uow_completeness_status', ( db_schema.apps.CompletenessStatusEnum.ACTIVE, db_schema.apps.CompletenessStatusEnum.MIN_COMPLETE, ) ) def test_optimized_uow_dependencies_not_deligate_if_active(uow_completeness_status, session): slz.ContentStatusFactory( unit_of_work=slz.UnitOfWorkFactory( completeness_status='ACTIVE', unit_of_work_code='spotify-20200101-sme-streams-v1', ), ) deps = unittest.mock.Mock(spec=EtlUoWDependencies) deps_optimized = UoWOptimizedDependencies( unittest.mock.Mock(), decorated=deps, session=session, dependencies=Dependencies([ EtlDependency( type='slz_uow', name='spotify-20200101-sme-streams-v1', is_optional=False ) ]), uow=apps.UnitOfWorkFactory( report_date=datetime.date(2020, 1, 1), completeness_status=uow_completeness_status, latest_content_status_timestamp=None, ), ) deps_optimized.load() assert not deps.load.called @pytest.mark.integration @pytest.mark.parametrize( 'gras_delta_1, gras_delta_2, can_deligate', ((SLZ_GRAS_UOWS_DELAY - 1, SLZ_GRAS_UOWS_DELAY - 1, False), (SLZ_GRAS_UOWS_DELAY - 1, SLZ_GRAS_UOWS_DELAY + 1, False), (SLZ_GRAS_UOWS_DELAY + 1, SLZ_GRAS_UOWS_DELAY + 1, True)) ) def test_gras_optimized_uow_dependencies_deligated( gras_delta_1, gras_delta_2, can_deligate, session ): gras_treshold_dt_1 = datetime.datetime.now(datetime.timezone.utc ) - datetime.timedelta(hours=gras_delta_1) gras_treshold_dt_2 = datetime.datetime.now(datetime.timezone.utc ) - datetime.timedelta(hours=gras_delta_2) slz.ContentStatusFactory( unit_of_work=slz.UnitOfWorkFactory( completeness_status='COMPLETE', unit_of_work_code='gras-20210219-sme-dim_participant-v1', last_updated_at=gras_treshold_dt_1 ), ) slz.ContentStatusFactory( unit_of_work=slz.UnitOfWorkFactory( completeness_status='COMPLETE', unit_of_work_code='gras-20210219-sme-dim_participant_member-v1', last_updated_at=gras_treshold_dt_2 ), ) deps = unittest.mock.Mock(spec=EtlUoWDependencies) deps_optimized = UoWOptimizedDependencies( unittest.mock.Mock(), decorated=deps, session=session, dependencies=Dependencies([ EtlDependency(type='slz_uow', name='gras-{yyyymmdd}-sme-dim_participant-v1'), EtlDependency(type='slz_uow', name='gras-{yyyymmdd}-sme-dim_participant_member-v1') ]), uow=apps.UnitOfWorkFactory( unit_of_work_code='gras-20210219-v1', report_date=datetime.date(2021, 2, 19), completeness_status=db_schema.apps.CompletenessStatusEnum.ACTIVE, latest_content_status_timestamp=None, ), ) deps_optimized.load() assert deps.load.called == can_deligate @pytest.mark.integration @pytest.mark.parametrize( 'completed_at,latest_ts,has_updates', ( ( datetime.datetime(2020, 1, 1, 10), datetime.datetime(2020, 1, 1, 12), False, ), ( datetime.datetime(2020, 1, 1, 12), datetime.datetime(2020, 1, 1, 10), True, ), ) ) def test_optimized_uow_dependencies_deligate_if_min_complete_with_statuses( completed_at, latest_ts, has_updates, session, ): slz.ContentStatusFactory( content_status='COMPLETE', completed_at=completed_at, unit_of_work=slz.UnitOfWorkFactory( completeness_status='MIN_COMPLETE', unit_of_work_code='spotify-20200101-sme-streams-v1', ), ) deps = unittest.mock.Mock(spec=EtlUoWDependencies) deps_optimized = UoWOptimizedDependencies( unittest.mock.Mock(), decorated=deps, session=session, dependencies=Dependencies([ EtlDependency( type='slz_uow', name='spotify-20200101-sme-streams-v1', ) ]), uow=apps.UnitOfWorkFactory( report_date=datetime.date(2020, 1, 1), completeness_status='ACTIVE', latest_content_status_timestamp=latest_ts, ), ) session.commit() deps_optimized.load() # if there are any updates - # we should deligate `.load` method to decorated instance assert has_updates is deps.load.called @pytest.mark.integration @pytest.mark.parametrize('is_force_complete,deligated', ( (False, False), (True, True), )) def test_optimized_uow_dependencies_loaded_for_force_completed( is_force_complete, deligated, session, ): slz_content_status_completed_at = datetime.datetime(2020, 1, 1, 10) apps_latest_content_status_timestamp = datetime.datetime(2020, 1, 1, 12) # slz_content_status_completed_at < apps_latest_content_status_timestamp - no updates slz.ContentStatusFactory( content_status='COMPLETE', completed_at=slz_content_status_completed_at, unit_of_work=slz.UnitOfWorkFactory( completeness_status='COMPLETE', unit_of_work_code='spotify-20200101-sme-streams-v1', is_force_complete=is_force_complete, ), ) deps = unittest.mock.Mock(spec=EtlUoWDependencies) deps_optimized = UoWOptimizedDependencies( unittest.mock.Mock(), decorated=deps, session=session, dependencies=Dependencies([ EtlDependency( type='slz_uow', name='spotify-20200101-sme-streams-v1', ) ]), uow=apps.UnitOfWorkFactory( report_date=datetime.date(2020, 1, 1), completeness_status='ACTIVE', latest_content_status_timestamp=apps_latest_content_status_timestamp, ), ) session.commit() deps_optimized.load() # if there are any force completed dependent slz uows - # we should deligate `.load` method to decorated instance assert deligated is deps.load.called @pytest.mark.integration @pytest.mark.parametrize('is_force_complete,deligated', ( (False, False), (True, True), )) def test_optimized_uow_dependencies_loaded_for_force_completed_optional_slz_deps( is_force_complete, deligated, session, ): slz_content_status_completed_at = datetime.datetime(2021, 10, 3, 10) apps_latest_content_status_timestamp = datetime.datetime(2021, 10, 3, 12) # slz_content_status_completed_at < apps_latest_content_status_timestamp - no updates slz.ContentStatusFactory( content_status='COMPLETE', completed_at=slz_content_status_completed_at, unit_of_work=slz.UnitOfWorkFactory( completeness_status='COMPLETE', unit_of_work_code='amazonmusicunlimited-20211003-theorchard-activity-v1', is_force_complete=is_force_complete, ), ) deps = unittest.mock.Mock(spec=EtlUoWDependencies) deps_optimized = UoWOptimizedDependencies( unittest.mock.Mock(), decorated=deps, session=session, dependencies=Dependencies([ EtlDependency( type='slz_uow', name='amazonmusicunlimited-20211003-theorchard-activity-v1', is_optional=True ) ]), uow=apps.UnitOfWorkFactory( report_date=datetime.date(2021, 10, 3), completeness_status='ACTIVE', latest_content_status_timestamp=apps_latest_content_status_timestamp, ), ) session.commit() deps_optimized.load() # if there are any force completed dependent optional slz uows - # we should deligate `.load` method to decorated instance assert deligated is deps.load.called @pytest.mark.integration @pytest.mark.parametrize( 'slz_uow_completeness_status, etl_uow_completeness_status, is_ready', ( ( db_schema.slz.CompletenessStatusEnum.COMPLETE, db_schema.apps.CompletenessStatusEnum.COMPLETE, True, ), ( db_schema.slz.CompletenessStatusEnum.MIN_COMPLETE, db_schema.apps.CompletenessStatusEnum.MIN_COMPLETE, True, ), ( db_schema.slz.CompletenessStatusEnum.COMPLETE, db_schema.apps.CompletenessStatusEnum.MIN_COMPLETE, True, ), ( db_schema.slz.CompletenessStatusEnum.MIN_COMPLETE, db_schema.apps.CompletenessStatusEnum.COMPLETE, True, ), ( db_schema.slz.CompletenessStatusEnum.COMPLETE, db_schema.apps.CompletenessStatusEnum.ACTIVE, False, ), ( db_schema.slz.CompletenessStatusEnum.COMPLETE, db_schema.apps.CompletenessStatusEnum.FAILED, False, ), ( db_schema.slz.CompletenessStatusEnum.COMPLETE, db_schema.apps.CompletenessStatusEnum.QUEUED, False, ), ) ) def test_combined_slz_and_etl_depdendencies_is_ready( slz_uow_completeness_status, etl_uow_completeness_status, is_ready, session, ): etl_uow = apps.UnitOfWorkFactory(report_date=datetime.date(2020, 1, 1)) # create dependent slz uow slz.UnitOfWorkFactory( unit_of_work_code='spotify-20200101-sme-aggregatedstreams-v2', report_date=datetime.date(2020, 1, 1), completeness_status=slz_uow_completeness_status, ) # create dependent etl uow apps.UnitOfWorkFactory( unit_of_work_code='spotify-20200101-dg_tracks', completeness_status=etl_uow_completeness_status, report_date=datetime.date(2020, 1, 1), ) dependencies = [ EtlDependency(type='slz_uow', name='spotify-{yyyymmdd}-sme-aggregatedstreams-v2'), EtlDependency(type='etl_uow', name='spotify-{yyyymmdd}-dg_tracks'), ] uow_dependencies = UoWOptimizedDependencies( decorated=EtlUoWDependencies( unittest.mock.Mock(), session, unittest.mock.Mock(), dependencies=Dependencies(dependencies), date=etl_uow.report_date, include_on_hold=False, dsp='spotify', env='test', ), logger=unittest.mock.Mock(), session=session, dependencies=dependencies, uow=etl_uow, ) uow_dependencies.load() assert uow_dependencies.is_ready() == is_ready if is_ready: # we have dependent_entities only for slz dependencies assert len(uow_dependencies.dependent_entities()) == 1 @pytest.mark.integration @pytest.mark.parametrize( 'slz_uow_completeness_status, is_snowflake_dep_updated, is_ready', ( ( db_schema.slz.CompletenessStatusEnum.COMPLETE, True, True, ), ( db_schema.slz.CompletenessStatusEnum.MIN_COMPLETE, True, True, ), ( db_schema.slz.CompletenessStatusEnum.COMPLETE, False, False, ), ( db_schema.slz.CompletenessStatusEnum.ACTIVE, True, False, ), ) ) @unittest.mock.patch('apps_etl_manager.dependencies.dependencies.SnowflakeDependency.load') def test_combined_slz_and_sf_depdendencies_is_ready( is_snowflake_dep_updated_mock, slz_uow_completeness_status, is_snowflake_dep_updated, is_ready, session, ): etl_uow = apps.UnitOfWorkFactory(report_date=datetime.date(2021, 2, 15)) # create dependent slz uow slz.UnitOfWorkFactory( unit_of_work_code='youtubereporting-20210215-sme-asset_combined-a2', report_date=datetime.date(2021, 2, 15), completeness_status=slz_uow_completeness_status, ) dependencies = [ EtlDependency(type='slz_uow', name='youtubereporting-{yyyymmdd}-sme-asset_combined-a2'), EtlDependency(type='snowflake', name='V_YT_VIDEO_ISRC'), ] uow_dependencies = UoWOptimizedDependencies( decorated=EtlUoWDependencies( unittest.mock.Mock(), session, unittest.mock.Mock(), dependencies=Dependencies(dependencies), date=etl_uow.report_date, include_on_hold=False, dsp='youtubereporting', env='test', ), logger=unittest.mock.Mock(), session=session, dependencies=dependencies, uow=etl_uow, ) is_snowflake_dep_updated_mock.return_value = is_snowflake_dep_updated uow_dependencies.load() assert uow_dependencies.is_ready() == is_ready if is_ready: # we have dependent_entities only for slz dependencies assert len(uow_dependencies.dependent_entities()) == 1 @pytest.mark.integration @pytest.mark.parametrize( 'users_completeness_status, dsp, dependent_entities_number', ( (db_schema.slz.CompletenessStatusEnum.COMPLETE, 'spotify', 2), (db_schema.slz.CompletenessStatusEnum.COMPLETE, 'not_spotify', 2), (db_schema.slz.CompletenessStatusEnum.ACTIVE, 'not_spotify', 2), (db_schema.slz.CompletenessStatusEnum.ACTIVE, 'spotify', 0), ) ) def test_depdendencies_load_excludes_joined_by_licensor_optional_deps( session, users_completeness_status, dsp, dependent_entities_number ): apps_uow = apps.UnitOfWorkFactory(report_date=datetime.date(2021, 11, 16)) slz.UnitOfWorkFactory( unit_of_work_code='spotify-20211116-theorchard-streams-v4', report_date=datetime.date(2021, 11, 16), completeness_status=db_schema.slz.CompletenessStatusEnum.COMPLETE, ) slz.UnitOfWorkFactory( unit_of_work_code='spotify-20211116-theorchard-users-v4', report_date=datetime.date(2021, 11, 16), completeness_status=users_completeness_status, ) uow_dependencies = EtlUoWDependencies( unittest.mock.Mock(), session, unittest.mock.Mock(), dependencies=Dependencies([ EtlDependency( type='slz_uow', name='spotify-{yyyymmdd}-theorchard-streams-v4', is_optional=True, ), EtlDependency( type='slz_uow', name='spotify-{yyyymmdd}-theorchard-users-v4', is_optional=True, ), ]), date=apps_uow.report_date, include_on_hold=False, dsp=dsp, env='test', ) session.commit() uow_dependencies.load() dependent_entities = uow_dependencies.dependent_entities() assert len(dependent_entities) == dependent_entities_number