import json import os import unittest import unittest.mock from contextlib import contextmanager import pytest from exp_manager_lambda import get_config_from_file from exp_manager_lambda.entities import Config from exp_manager_lambda.handler import handler FIXTURES_PATH = os.path.join(os.path.dirname(__file__), "fixtures") @contextmanager def ctx_manager_mock(**kwargs): yield int @pytest.mark.parametrize( 'uow_id, decompressed_paths, expected_calls_to_dcs_table', [ # regular message from slz flow, separate record for each file ( 'spotify-20170829-smejp-users-v1', 's3://dev-archive/spotify/path1.txt,s3://dev-archive/spotify/path2.txt', 2, ), # regular message from slz flow, but we expect parquets for this report ( 'youtubereporting-20201108-sme-content_owner_asset-a2', 's3://dev-archive/spotify/path1.txt', 0, ), # message from etl flow, we expect parquets for this report. One record for all parquets ( 'youtubereporting-20201108-sme-content_owner_asset-a2', 's3://dev-archive/spotify/part-00000-e408e7f5-2f61-4bc4-8c9b-d4dd6a6d600f.c000.snappy.parquet,' 's3://dev-archive/spotify/part-00001-e408e7f5-2f61-4bc4-8c9b-d4dd6a6d600f.c000.snappy.parquet', 1, ), # message from etl flow, we expect parquets for this report. One record for all parquets ( 'youtubereporting-20201108-sme-content_owner_asset-a2', 's3://dev-archive/spotify/part-00000-e408e7f5-2f61-4bc4-8c9b-d4dd6a6d600f.c000.snappy.parquet', 1, ), # message from etl flow, but we process record in main flow using slz files for this report. # should be skipped ( 'spotify-20170829-smejp-users-v1', 's3://dev-archive/spotify/part-00000-e408e7f5-2f61-4bc4-8c9b-d4dd6a6d600f.c000.snappy.parquet,' 's3://dev-archive/spotify/part-00001-e408e7f5-2f61-4bc4-8c9b-d4dd6a6d600f.c000.snappy.parquet', 0, ), ] ) @unittest.mock.patch('boto3.client') def test_handler_ok( mocked_client, uow_id, decompressed_paths, expected_calls_to_dcs_table ): event_base = { 'ContentName': 'content-name', 'Context': 'context', 'CompressedPath': 'path', 'OptionalConfig': {'RDS_SECRET_KEY': 'test/key'} } repo = unittest.mock.Mock() config = Config( environment='dev', rds_secrets_key='/delphi/dev/key', trace_id='123', excluded_dsp=['apple'], exp_dsp_config=get_config_from_file( os.path.join(FIXTURES_PATH, "exploration-dsp-config_full.json") ), ) repo.advisory_locked_transaction.side_effect = ctx_manager_mock mocks = { 'rds': unittest.mock.Mock(), 'secretsmanager': unittest.mock.Mock(), } mocked_client.side_effect = lambda res: mocks[res] event = {'UoWID': uow_id, 'DecompressedPaths': decompressed_paths, **event_base} payload = {'messageId': 'sqs-message-id', 'body': json.dumps(event)} result = handler( unittest.mock.Mock(), payload, repo, unittest.mock.Mock(), config ) assert repo.create_disassemble_content_status_by_uow_id.call_count == expected_calls_to_dcs_table assert result is True