from unittest import mock import pytest from dapd_scheduler_lambda import config, const from dapd_scheduler_lambda.repository import batch @pytest.fixture(scope='function') def batch_repository_params(): return config.Batch( queues=config.Queues( low={ 'JobQueue': 'job-queue/dev-delphi-low', 'JobDefinition': 'job-definition/dev-delphi-dapd-ingestion:1', }, medium={ 'JobQueue': 'job-queue/dev-delphi-medium', 'JobDefinition': 'job-definition/dev-delphi-dapd-ingestion:2', }, high={ 'JobQueue': 'job-queue/dev-delphi-high', 'JobDefinition': 'job-definition/dev-delphi-dapd-ingestion:3', }, ), command=config.Command( tasks_cleanup_on_success=True, results_bucket='results_bucket', corrupted_bucket='corrupted_bucket', validation_schema_s3_path='s3://bucket/public-data/schemas', workflow_db_secret_key='delphi/test/pg/workflowdb', credentials={ const.DataSourceEnum.SPOTIFY: config.DataSourceCredentials( sme_secret_key='sme/spotify/key', user_secret_key='user/spotify/key', ), const.DataSourceEnum.APPLE_MUSIC: config.DataSourceCredentials(sme_secret_key='sme/apple-music/key'), }, batch_sync=False, batch_concurrency=1, ) ) @pytest.mark.parametrize( 'priority,expected', [ ( 1, { 'JobQueue': 'job-queue/dev-delphi-high', 'JobDefinition': 'job-definition/dev-delphi-dapd-ingestion:3', }, ), ( 4, { 'JobQueue': 'job-queue/dev-delphi-medium', 'JobDefinition': 'job-definition/dev-delphi-dapd-ingestion:2', }, ), ( 5, { 'JobQueue': 'job-queue/dev-delphi-medium', 'JobDefinition': 'job-definition/dev-delphi-dapd-ingestion:2', }, ), ( 7, { 'JobQueue': 'job-queue/dev-delphi-low', 'JobDefinition': 'job-definition/dev-delphi-dapd-ingestion:1', }, ), ] ) def test_batch_config(priority, expected, batch_repository_params): logger = mock.Mock() batch_client = mock.Mock() repo = batch.Batch(logger, 'test', batch_repository_params, batch_client) result = repo.get_config(priority) assert result == expected @pytest.mark.parametrize( 'dsp,entity,tasks_cleanup_on_success,is_removed_ttl,expected', [ ( const.DataSourceEnum.SPOTIFY, const.DataSourceEntityEnum.PLAYLIST, True, '1_DAY', [ '--env test', '--data-source-name spotify', '--data-source-entity playlist', '--data-source-sme-secret-key sme/spotify/key', '--workflow-db-secret-key delphi/test/pg/workflowdb', '--kinesis-transform-stream test-delphi-dapd-spotify-playlists', '--s3-backup-bucket results_bucket', '--corrupted-bucket corrupted_bucket', '--tasks-file s3://bucket/tasks/123.json', '--validation-schema-storage s3://bucket/public-data/schemas', '--sentry-secret-key delphi/dev/sentry_config', '--is-removed-ttl 1_DAY', '--concurrency 1', '--data-source-user-secret-key user/spotify/key', '--tasks-file-delete', ], ), ( const.DataSourceEnum.SPOTIFY, const.DataSourceEntityEnum.TRACK, False, '30_DAYS', [ '--env test', '--data-source-name spotify', '--data-source-entity track', '--data-source-sme-secret-key sme/spotify/key', '--workflow-db-secret-key delphi/test/pg/workflowdb', '--kinesis-transform-stream test-delphi-dapd-spotify-tracks', '--s3-backup-bucket results_bucket', '--corrupted-bucket corrupted_bucket', '--tasks-file s3://bucket/tasks/123.json', '--validation-schema-storage s3://bucket/public-data/schemas', '--sentry-secret-key delphi/dev/sentry_config', '--is-removed-ttl 30_DAYS', '--concurrency 1', '--data-source-user-secret-key user/spotify/key', ] ), ( const.DataSourceEnum.APPLE_MUSIC, const.DataSourceEntityEnum.ALBUM, True, 'ASkjshad1', [ '--env test', '--data-source-name apple_music', '--data-source-entity album', '--data-source-sme-secret-key sme/apple-music/key', '--workflow-db-secret-key delphi/test/pg/workflowdb', '--kinesis-transform-stream test-delphi-dapd-apple_music-albums', '--s3-backup-bucket results_bucket', '--corrupted-bucket corrupted_bucket', '--tasks-file s3://bucket/tasks/123.json', '--validation-schema-storage s3://bucket/public-data/schemas', '--sentry-secret-key delphi/dev/sentry_config', '--is-removed-ttl ASkjshad1', '--concurrency 1', '--tasks-file-delete', ] ), ( const.DataSourceEnum.APPLE_MUSIC, const.DataSourceEntityEnum.ARTIST, False, '123', [ '--env test', '--data-source-name apple_music', '--data-source-entity artist', '--data-source-sme-secret-key sme/apple-music/key', '--workflow-db-secret-key delphi/test/pg/workflowdb', '--kinesis-transform-stream test-delphi-dapd-apple_music-artists', '--s3-backup-bucket results_bucket', '--corrupted-bucket corrupted_bucket', '--tasks-file s3://bucket/tasks/123.json', '--validation-schema-storage s3://bucket/public-data/schemas', '--sentry-secret-key delphi/dev/sentry_config', '--is-removed-ttl 123', '--concurrency 1', ] ), ] ) def test_batch_get_command( dsp, entity, tasks_cleanup_on_success, is_removed_ttl, expected, batch_repository_params ): logger = mock.Mock() batch_client = mock.Mock() batch_repository_params.command.tasks_cleanup_on_success = tasks_cleanup_on_success repo = batch.Batch(logger, 'test', batch_repository_params, batch_client) task_path = 's3://bucket/tasks/123.json' creds = batch_repository_params.command.credentials[dsp] result = repo.get_command( dsp.value, entity.value, task_path=task_path, sentry_secret_key='delphi/dev/sentry_config', sme_secret_key=creds.sme_secret_key, user_secret_key=creds.user_secret_key, is_removed_ttl=is_removed_ttl, batch_sync=False, batch_concurrency=1, ) assert result == expected @pytest.mark.parametrize( 'dsp,name,submit_job_args,expected', [ ( const.DataSourceEnum.SPOTIFY, const.DataSourceEntityEnum.PLAYLIST, dict( containerOverrides={ 'command': [ '--env test', '--data-source-name spotify', '--data-source-entity playlist', '--data-source-sme-secret-key sme/spotify/key', '--workflow-db-secret-key delphi/test/pg/workflowdb', '--kinesis-transform-stream test-delphi-dapd-spotify-playlists', '--s3-backup-bucket results_bucket', '--corrupted-bucket corrupted_bucket', '--tasks-file s3://bucket/task_1.json', '--validation-schema-storage s3://bucket/public-data/schemas', '--sentry-secret-key delphi/dev/sentry_config', '--is-removed-ttl 5_DAYS', '--concurrency 1', '--data-source-user-secret-key user/spotify/key', '--tasks-file-delete', ], 'memory': 7000, 'vcpus': 2, }, jobDefinition='job-definition/dev-delphi-dapd-ingestion:2', jobName='spotify_playlist_job_id', jobQueue='job-queue/dev-delphi-medium', timeout={'attemptDurationSeconds': 30 * 60}, retryStrategy={ 'attempts': 3, "evaluateOnExit": [ { "onExitCode": "137", "action": "RETRY", }, { "onExitCode": "0", "action": "EXIT", }, ] }, ), True, ), ( const.DataSourceEnum.APPLE_MUSIC, const.DataSourceEntityEnum.PLAYLIST, dict( containerOverrides={ 'command': [ '--env test', '--data-source-name apple_music', '--data-source-entity playlist', '--data-source-sme-secret-key sme/apple-music/key', '--workflow-db-secret-key delphi/test/pg/workflowdb', '--kinesis-transform-stream test-delphi-dapd-apple_music-playlists', '--s3-backup-bucket results_bucket', '--corrupted-bucket corrupted_bucket', '--tasks-file s3://bucket/task_1.json', '--validation-schema-storage s3://bucket/public-data/schemas', '--sentry-secret-key delphi/dev/sentry_config', '--is-removed-ttl 5_DAYS', '--concurrency 1', '--tasks-file-delete', ], 'memory': 7000, 'vcpus': 2, }, jobDefinition='job-definition/dev-delphi-dapd-ingestion:2', jobName='apple_music_playlist_job_id', jobQueue='job-queue/dev-delphi-medium', timeout={'attemptDurationSeconds': 30 * 60}, retryStrategy={ 'attempts': 3, "evaluateOnExit": [ { "onExitCode": "137", "action": "RETRY", }, { "onExitCode": "0", "action": "EXIT", }, ] }, ), True, ), ] ) def test_batch_submit_task(dsp, name, submit_job_args, expected, batch_repository_params): logger = mock.Mock() batch_client = mock.Mock() repo = batch.Batch(logger, 'test', batch_repository_params, batch_client) task_path = 's3://bucket/task_1.json' priority = 5 credentials = batch_repository_params.command.credentials[dsp] result = repo.submit_task( 'job_id', dsp.value, name.value, task_path, priority, 'delphi/dev/sentry_config', credentials.sme_secret_key, '5_DAYS', user_secret_key=credentials.user_secret_key, batch_sync=False, batch_concurrency=1, ) assert result == expected batch_client.submit_job.assert_called_with(**submit_job_args)