#!/usr/bin/env python """Lambda test module.""" import json import os from unittest.mock import call from unittest.mock import MagicMock from aws_kinesis_agg import deaggregator import boto3 import pytest import config # noqa import const # noqa import index # noqa from logic import artist_info from logic import product from logic import project from logic import subaccount from logic import track from logic import users from logic import vendor from models import ows_account from models import ows_artist from models import ows_product from models import ows_project_manager from models import ows_track from models import ows_users def test_handler_process_project( mocker, project_kinesis_event_fixture, project_cloudsearch_document, project_document_fixture, converted_project_event_fixture): """Test index.handler push a project event to Kinesis.""" mocker.patch.object( ows_project_manager, 'get_project_document', return_value=project_document_fixture) mocker.spy(project, 'prepare_doc') mocker.spy(deaggregator, 'iter_deaggregate_records') client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) index.handler(project_kinesis_event_fixture, None) project.prepare_doc.assert_called_with(converted_project_event_fixture) assert boto3.client.call_args[1].get( 'endpoint_url') == os.environ.get('CLOUDSEARCH_ENDPOINT_PROJECTS') deaggregator.iter_deaggregate_records.assert_called_with( project_kinesis_event_fixture['Records']) client_mock.upload_documents.assert_called() upload_documents_args = client_mock.upload_documents.call_args[1] documents_arg = json.loads(upload_documents_args['documents']) assert documents_arg == [project_cloudsearch_document] def test_handler_process_release_approval_queue( mocker, release_approval_queue_kinesis_event_fixture, product_cloudsearch_document, product_document_fixture, converted_release_approval_queue_event_fixture): """Test index.handler process release_approval_queue event from Kinesis.""" _test_handler_process_product( mocker, release_approval_queue_kinesis_event_fixture, product_cloudsearch_document, product_document_fixture, converted_release_approval_queue_event_fixture) def test_handler_process_release_correction( mocker, release_correction_kinesis_event_fixture, product_cloudsearch_document, product_document_fixture, converted_release_correction_event_fixture): """Test index.handler process release_correction event from Kinesis.""" _test_handler_process_product( mocker, release_correction_kinesis_event_fixture, product_cloudsearch_document, product_document_fixture, converted_release_correction_event_fixture) def test_handler_process_product( mocker, product_kinesis_event_fixture, product_cloudsearch_document, product_document_fixture, converted_product_event_fixture): """Test index.handler process product event from Kinesis.""" _test_handler_process_product( mocker, product_kinesis_event_fixture, product_cloudsearch_document, product_document_fixture, converted_product_event_fixture) def test_handler_process_update_product( mocker, product_kinesis_update_event_fixture, product_cloudsearch_document, product_document_fixture, converted_product_update_event_fixture): """Test index.handler process product update event from Kinesis.""" client_mock = MagicMock() aws_client_mock = mocker.patch.object(boto3, 'client', return_value=client_mock) mocker.spy(deaggregator, 'iter_deaggregate_records') release_doc = { 'type': 'add', 'id': 'rel1', 'fields': {'release_id': 1}} tracks_docs = [{'track_unique_id': 123}, {'track_unique_id': 456}] product_prepare_doc_mock = mocker.patch('logic.product.prepare_doc', return_value=release_doc) track_prepare_release_tracks_docs_mock = mocker.patch( 'logic.track.prepare_release_tracks_docs', return_value=tracks_docs) index.handler(product_kinesis_update_event_fixture, None) product_prepare_doc_mock.assert_called_with(converted_product_update_event_fixture) track_prepare_release_tracks_docs_mock.assert_called_with(converted_product_update_event_fixture) deaggregator.iter_deaggregate_records.assert_called_with(product_kinesis_update_event_fixture['Records']) aws_client_mock.call_args.call_count == 2 aws_client_mock.call_args.mock_calls == [ call('cloudsearch', endpoint_url=os.environ.get('CLOUDSEARCH_ENDPOINT_RELEASES')), call('cloudsearch', endpoint_url=os.environ.get('CLOUDSEARCH_ENDPOINT')) ] client_mock.upload_documents.call_count == 2 client_mock.upload_documents.mock_calls == [ call(documents='[{"release_id": 1}]', contentType='application/json'), call(documents='[{"track_unique_id": 123}, {"track_unique_id": 456}]', contentType='application/json') ] def _test_handler_process_product( mocker, kinesis_event, product_cloudsearch_document, product_document_fixture, converted_event): """Test index.handler push a product event to Kinesis.""" mocker.patch.object( ows_product, 'get_product_document', return_value=product_document_fixture) mocker.spy(product, 'prepare_doc') mocker.spy(deaggregator, 'iter_deaggregate_records') client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) index.handler(kinesis_event, None) product.prepare_doc.assert_called_with(converted_event) assert boto3.client.call_args[1].get( 'endpoint_url') == os.environ.get('CLOUDSEARCH_ENDPOINT_RELEASES') deaggregator.iter_deaggregate_records.assert_called_with( kinesis_event['Records']) client_mock.upload_documents.assert_called() upload_documents_args = client_mock.upload_documents.call_args[1] documents_arg = json.loads(upload_documents_args['documents']) assert documents_arg == [product_cloudsearch_document] def test_handler_process_track( mocker, converted_track_event_fixture, track_kinesis_event_fixture, track_cloudsearch_document, get_release_info_fixture, get_track_info_fixture): """Test index.handler push a track event to Kinesis.""" mocker.patch.object( ows_product, 'get_release_info', return_value=get_release_info_fixture) mocker.patch.object( ows_track, 'get_track_artist_info', return_value=get_track_info_fixture) mocker.spy(track, 'prepare_doc') mocker.spy(deaggregator, 'iter_deaggregate_records') client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) index.handler(track_kinesis_event_fixture, None) assert boto3.client.call_args[1].get( 'endpoint_url') == os.environ.get('CLOUDSEARCH_ENDPOINT') track.prepare_doc.assert_called_with(converted_track_event_fixture) deaggregator.iter_deaggregate_records.assert_called_with( track_kinesis_event_fixture['Records']) client_mock.upload_documents.assert_called() upload_documents_args = client_mock.upload_documents.call_args[1] documents_arg = json.loads(upload_documents_args['documents']) assert documents_arg == [track_cloudsearch_document] def test_handler_exception_track( mocker, converted_track_event_fixture, track_kinesis_event_fixture, get_release_info_fixture): """Test index.handler track exception skips upload.""" mocker.patch.object( ows_product, 'get_release_info', return_value=get_release_info_fixture) mocker.patch.object( ows_track, 'get_track_artist_info', side_effect=Exception(const.OWS_TRACK_ERROR)) mocker.spy(track, 'prepare_doc') client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) with pytest.raises(Exception, match=const.OWS_TRACK_ERROR): index.handler(track_kinesis_event_fixture, None) track.prepare_doc.assert_called_with(converted_track_event_fixture) assert not client_mock.upload_documents.called def test_handler_process_artist( mocker, artist_kinesis_event_fixture, artist_cloudsearch_document, artist_document_fixture, converted_artist_event_fixture): """Test index.handler push an artist event to Kinesis.""" mocker.patch.object( ows_artist, 'get_artist_document', return_value=artist_document_fixture) mocker.spy(artist_info, 'prepare_doc') mocker.spy(deaggregator, 'iter_deaggregate_records') client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) index.handler(artist_kinesis_event_fixture, None) artist_info.prepare_doc.assert_called_with(converted_artist_event_fixture) assert boto3.client.call_args[1].get( 'endpoint_url') == os.environ.get('CLOUDSEARCH_ENDPOINT_ARTISTS') deaggregator.iter_deaggregate_records.assert_called_with( artist_kinesis_event_fixture['Records']) client_mock.upload_documents.assert_called() upload_documents_args = client_mock.upload_documents.call_args[1] documents_arg = json.loads(upload_documents_args['documents']) assert documents_arg == [artist_cloudsearch_document] def test_handler_exception_project( mocker, converted_project_event_fixture, project_kinesis_event_fixture, project_document_fixture): """Test index.handler project exception skips upload.""" mocker.patch.object( ows_project_manager, 'get_project_document', side_effect=Exception(const.OWS_PROJECT_ERROR)) mocker.spy(project, 'prepare_doc') client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) with pytest.raises(Exception, match=const.OWS_PROJECT_ERROR): index.handler(project_kinesis_event_fixture, None) project.prepare_doc.assert_called_with(converted_project_event_fixture) assert not client_mock.upload_documents.called def test_handler_exception_artists( mocker, converted_artist_event_fixture, artist_kinesis_event_fixture, artist_document_fixture): """Test index.handler artist exception skips upload.""" mocker.patch.object( ows_artist, 'get_artist_document', side_effect=Exception(const.OWS_ARTIST_ERROR)) mocker.spy(artist_info, 'prepare_doc') client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) with pytest.raises(Exception, match=const.OWS_ARTIST_ERROR): index.handler(artist_kinesis_event_fixture, None) artist_info.prepare_doc.assert_called_with(converted_artist_event_fixture) assert not client_mock.upload_documents.called def test_handler_exception_product( mocker, converted_product_event_fixture, product_kinesis_event_fixture, product_document_fixture): """Test index.handler product exception skips upload.""" mocker.patch.object( ows_product, 'get_product_document', side_effect=Exception(const.OWS_PRODUCT_ERROR)) mocker.spy(product, 'prepare_doc') client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) with pytest.raises(Exception, match=const.OWS_PRODUCT_ERROR): index.handler(product_kinesis_event_fixture, None) product.prepare_doc.assert_called_with(converted_product_event_fixture) assert not client_mock.upload_documents.called def test_handler_incompatible_doc(mocker, kinesis_event_incompatible_fixture): """Test index.handler don't upload empty document batches.""" client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) index.handler(kinesis_event_incompatible_fixture, None) client_mock.upload_documents.assert_not_called() def test_handler_cloudsearch_upload_exception( mocker, product_kinesis_event_fixture, product_cloudsearch_document, product_document_fixture, converted_product_event_fixture): """Test Cloudsearch exception logs to Sentry.""" mocker.patch.object( ows_product, 'get_product_document', return_value=product_document_fixture) mocker.spy(product, 'prepare_doc') mocker.spy(deaggregator, 'iter_deaggregate_records') client_mock = MagicMock() client_mock.upload_documents = MagicMock(side_effect=Exception()) mocker.patch.object( boto3, 'client', return_value=client_mock) with pytest.raises(Exception): index.handler(product_kinesis_event_fixture, None) product.prepare_doc.assert_called_with(converted_product_event_fixture) assert boto3.client.call_args[1].get( 'endpoint_url') == os.environ.get('CLOUDSEARCH_ENDPOINT_RELEASES') deaggregator.iter_deaggregate_records.assert_called_with( product_kinesis_event_fixture['Records']) client_mock.upload_documents.assert_called() upload_documents_args = client_mock.upload_documents.call_args[1] documents_arg = json.loads(upload_documents_args['documents']) assert documents_arg == [product_cloudsearch_document] def test_handler_process_subaccount( mocker, subaccount_kinesis_event_fixture, subaccount_cloudsearch_document, subaccount_document_fixture, converted_subaccount_event_fixture, user_document_fixture, user_cloudsearch_document_for_subaccount): """Test index.handler push a subaccount event to Kinesis.""" mocker.patch.object( ows_account, 'get_account_document', return_value=subaccount_document_fixture) mocker.patch.object( ows_users, 'get_user_document', return_value=user_document_fixture) mocker.spy(subaccount, 'prepare_doc') mocker.spy(users, 'prepare_docs') mocker.spy(deaggregator, 'iter_deaggregate_records') client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) index.handler(subaccount_kinesis_event_fixture, None) subaccount.prepare_doc.assert_called_with(converted_subaccount_event_fixture) users.prepare_docs.assert_called_with(converted_subaccount_event_fixture) assert boto3.client.call_args[1].get( 'endpoint_url') == os.environ.get('CLOUDSEARCH_ENDPOINT_LABELS') assert boto3.client.call_args[1].get( 'endpoint_url') == os.environ.get('CLOUDSEARCH_ENDPOINT_USERS') deaggregator.iter_deaggregate_records.assert_called_with( subaccount_kinesis_event_fixture['Records']) client_mock.upload_documents.assert_called() upload_documents_args = client_mock.upload_documents.call_args[1] documents_arg = json.loads(upload_documents_args['documents']) assert documents_arg == user_cloudsearch_document_for_subaccount def test_handler_process_vendor( mocker, vendor_kinesis_event_fixture, vendor_cloudsearch_document, vendor_document_fixture, converted_vendor_event_fixture): """Test index.handler push a vendor event to Kinesis.""" mocker.patch.object( ows_account, 'get_account_document', return_value=vendor_document_fixture) mocker.spy(vendor, 'prepare_doc') mocker.spy(deaggregator, 'iter_deaggregate_records') client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) index.handler(vendor_kinesis_event_fixture, None) vendor.prepare_doc.assert_called_with(converted_vendor_event_fixture) assert boto3.client.call_args[1].get( 'endpoint_url') == os.environ.get('CLOUDSEARCH_ENDPOINT_LABELS') deaggregator.iter_deaggregate_records.assert_called_with( vendor_kinesis_event_fixture['Records']) client_mock.upload_documents.assert_called() upload_documents_args = client_mock.upload_documents.call_args[1] documents_arg = json.loads(upload_documents_args['documents']) assert documents_arg == [vendor_cloudsearch_document] def test_handler_process_user_vend_contact( mocker, user_kinesis_event_fixture, user_cloudsearch_document, user_document_fixture, vendor_document_fixture, converted_vend_contact_event_fixture): """Test index.handler push a vend_contact event to Kinesis.""" document = user_document_fixture document.pop() mocker.patch.object( ows_users, 'get_user_document', return_value=document) mocker.patch.object( ows_account, 'get_account_document', return_value=vendor_document_fixture) mocker.spy(users, 'prepare_docs') mocker.spy(deaggregator, 'iter_deaggregate_records') client_mock = MagicMock() mocker.patch.object( boto3, 'client', return_value=client_mock) index.handler(user_kinesis_event_fixture, None) users.prepare_docs.assert_called_with(converted_vend_contact_event_fixture) assert boto3.client.call_args[1].get( 'endpoint_url') == os.environ.get('CLOUDSEARCH_ENDPOINT_USERS') deaggregator.iter_deaggregate_records.assert_called_with( user_kinesis_event_fixture['Records']) client_mock.upload_documents.assert_called() upload_documents_args = client_mock.upload_documents.call_args[1] documents_arg = json.loads(upload_documents_args['documents']) assert documents_arg == user_cloudsearch_document # it wont write to kafka for vend_contact domain