# pylint: disable=unused-argument,too-many-arguments,too-many-locals,protected-access import json from datetime import datetime, timedelta import pytest from dapd_transformation_service.entities.worker_control import ConsumerMetrics from dapd_transformation_service.helpers.misc import utcnow from dapd_transformation_service.services.consumer import ConsumerService pytestmark = pytest.mark.integration @pytest.mark.skip('Temporary skip') @pytest.mark.parametrize('entity', ['playlists', 'artists', 'albums', 'tracks']) def test_match_entity_with_stream( config, kinesis, db, mocker, dim_dsp, dim_playlist, dim_album, dim_artist, dim_track, entity, logger, ): stream_name = f'{config.env}-delphi-dapd-{dim_dsp.dsp_name}-{entity}' config.kinesis_stream_name = stream_name event_mocked = mocker.patch('threading.Event', autospec=True) store_data_mocked = mocker.patch( 'dapd_transformation_service.services.consumer.ConsumerService.store_data' ) expected_data = json.dumps({'entity': entity}).encode() consumer_metrics = ConsumerMetrics() upserter_service = mocker.MagicMock() service = ConsumerService(logger, config, db, consumer_metrics, upserter_service) service.shard_iterator_ttl = 1 service.run(thread_sync_event=event_mocked) store_data_mocked.assert_called_once_with(expected_data) assert consumer_metrics.consumed_records == 1 @pytest.mark.skip('Temporary skip') @pytest.mark.parametrize('entity', ['playlists', 'artists', 'albums', 'tracks']) def test_no_data_since_last_update( config, kinesis, db, mocker, dim_dsp, dim_playlist, dim_album, dim_artist, dim_track, entity, logger, ): stream_name = f'{config.env}-delphi-dapd-{dim_dsp.dsp_name}-{entity}' config.kinesis_stream_name = stream_name event_mocked = mocker.patch('threading.Event', autospec=True) store_data_mocked = mocker.patch( 'dapd_transformation_service.services.consumer.ConsumerService.store_data' ) for model in [dim_track, dim_album, dim_artist, dim_playlist]: model.updated_at = datetime.utcnow() + timedelta(days=1) model.is_stream_synced = True db.flush() consumer_metrics = ConsumerMetrics() upserter_service = mocker.MagicMock() service = ConsumerService(logger, config, db, consumer_metrics, upserter_service) service.shard_iterator_ttl = 1 service.run(thread_sync_event=event_mocked) assert not store_data_mocked.called assert consumer_metrics.consumed_records == 0 @pytest.mark.skip('Temporary skip') @pytest.mark.parametrize('entity', ['playlists', 'artists', 'albums', 'tracks']) def test_updated_at_is_none( config, kinesis, db, mocker, dim_dsp, dim_playlist, dim_album, dim_artist, dim_track, entity, logger, ): stream_name = f'{config.env}-delphi-dapd-{dim_dsp.dsp_name}-{entity}' config.kinesis_stream_name = stream_name event_mocked = mocker.patch('threading.Event', autospec=True) store_data_mocked = mocker.patch( 'dapd_transformation_service.services.consumer.ConsumerService.store_data' ) for model in [dim_track, dim_album, dim_artist, dim_playlist]: model.updated_at = None db.flush() expected_data = json.dumps({'entity': entity}).encode() consumer_metrics = ConsumerMetrics() upserter_service = mocker.MagicMock() service = ConsumerService(logger, config, db, consumer_metrics, upserter_service) service.shard_iterator_ttl = 1 service.run(thread_sync_event=event_mocked) store_data_mocked.assert_called_once_with(expected_data) assert consumer_metrics.consumed_records == 1 @pytest.mark.skip('Temporary skip') @pytest.mark.parametrize('entity', ['playlists', 'artists', 'albums', 'tracks']) def test_empty_query_for_instance(config, kinesis, db, mocker, dim_dsp, entity, logger): stream_name = f'{config.env}-delphi-dapd-{dim_dsp.dsp_name}-{entity}' config.kinesis_stream_name = stream_name event_mocked = mocker.patch('threading.Event', autospec=True) store_data_mocked = mocker.patch( 'dapd_transformation_service.services.consumer.ConsumerService.store_data' ) expected_data = json.dumps({'entity': entity}).encode() consumer_metrics = ConsumerMetrics() upserter_service = mocker.MagicMock() service = ConsumerService(logger, config, db, consumer_metrics, upserter_service) service.shard_iterator_ttl = 1 service.run(thread_sync_event=event_mocked) store_data_mocked.assert_called_once_with(expected_data) assert consumer_metrics.consumed_records == 1 @pytest.mark.skip('Temporary skip') @pytest.mark.parametrize('entity', ['playlists', 'artists', 'albums', 'tracks']) def test_get_dimension_last_update__ok( config, db, mocker, dim_dsp, dim_track, dim_album, dim_playlist, dim_artist, entity, logger, ): expected = utcnow() for model in [dim_track, dim_album, dim_artist, dim_playlist]: model.updated_at = expected model.is_stream_synced = True db.flush() entity = 'tracks' stream_name = f'{config.env}-delphi-dapd-{dim_dsp.dsp_name}-{entity}' config.kinesis_stream_name = stream_name consumer_metrics = ConsumerMetrics() upserter_service = mocker.MagicMock() service = ConsumerService(logger, config, db, consumer_metrics, upserter_service) actual = service._get_dimension_last_update() assert actual == expected @pytest.mark.skip('Temporary skip') @pytest.mark.parametrize('entity', ['playlists', 'artists', 'albums', 'tracks']) def test_get_dimension_last_update__stream_not_synced( config, db, mocker, dim_dsp, dim_track, dim_album, dim_playlist, dim_artist, entity, logger, ): expected = utcnow() for model in [dim_track, dim_album, dim_artist, dim_playlist]: model.updated_at = expected model.is_stream_synced = False db.flush() entity = 'tracks' stream_name = f'{config.env}-delphi-dapd-{dim_dsp.dsp_name}-{entity}' config.kinesis_stream_name = stream_name consumer_metrics = ConsumerMetrics() upserter_service = mocker.MagicMock() service = ConsumerService(logger, config, db, consumer_metrics, upserter_service) actual = service._get_dimension_last_update() assert actual is None @pytest.mark.skip('Temporary skip') @pytest.mark.parametrize('entity', ['playlists', 'artists', 'albums', 'tracks']) def test_get_dimension_last_update__updated_at_is_none( config, db, mocker, dim_dsp, dim_track, dim_album, dim_playlist, dim_artist, entity, logger, ): for model in [dim_track, dim_album, dim_artist, dim_playlist]: model.updated_at = None model.is_stream_synced = True db.flush() entity = 'tracks' stream_name = f'{config.env}-delphi-dapd-{dim_dsp.dsp_name}-{entity}' config.kinesis_stream_name = stream_name consumer_metrics = ConsumerMetrics() upserter_service = mocker.MagicMock() service = ConsumerService(logger, config, db, consumer_metrics, upserter_service) actual = service._get_dimension_last_update() assert actual is None