from decimal import Decimal from unittest.mock import ANY from unittest.mock import call from unittest.mock import Mock from unittest.mock import patch import uuid from botocore.exceptions import ClientError from celery.exceptions import MaxRetriesExceededError from moto import mock_aws import boto3 from oto import response from owsrequest.utils import mock_request import pytest from flexmock import flexmock from masters_registry import config from masters_registry.connectors import s3 from masters_registry.connectors import sentry from masters_registry.connectors import sqs from masters_registry.constant import api_const from masters_registry.constant import bulk_tasks_const from masters_registry.constant import error from masters_registry.constant import field_const from masters_registry.constant import opcode_const from masters_registry.logic import bulk_tasks as bulk_tasks_logic from masters_registry.logic import dms_carveout from masters_registry.logic import locks from masters_registry.logic import masters_registry from masters_registry.logic import ownership as ownership_logic from masters_registry.models import bulk_tasks from masters_registry.models import ownership from masters_registry.models import ows_conflict_manager from masters_registry.models import yt_ownership from masters_registry.tasks import bulk from tests.helpers import patches from tests.tasks import bulk_import_upcs_fixtures as bulk_import_fixtures from tests.test_utils import equal_dicts from masters_registry.connectors.dynamodb import client as outside_dynamodb_client from masters_registry.connectors.s3 import client as outside_s3_client dynamodb_resource = boto3.resource("dynamodb") s3_resource = boto3.resource("s3") sqs_resource = sqs.sqs_resource sqs_client = boto3.client('sqs') @pytest.fixture def client_error_fixture(): return ClientError({'Error': {}}, 'error') def client_error(): return ClientError({'Error': {}}, 'error') @pytest.fixture def sentry_captureException(monkeypatch): sentry_client = Mock() captureException = Mock() sentry_client.captureException = captureException monkeypatch.setattr(sentry, 'sentry_client', sentry_client) return captureException @pytest.fixture def sentry_captureMessage(monkeypatch): sentry_client = Mock() captureMessage = Mock() sentry_client.captureMessage = captureMessage monkeypatch.setattr(sentry, 'sentry_client', sentry_client) return captureMessage @pytest.fixture() def generate_reports(monkeypatch): generate_reports = Mock() monkeypatch.setattr(bulk, 'generate_reports', generate_reports) return generate_reports @pytest.fixture() def generate_report_logic(monkeypatch): generate_report = Mock() monkeypatch.setattr(bulk_tasks_logic, 'generate_report', generate_report) return generate_report def mock_retry(monkeypatch, bulk_task): retry = Mock() retry.side_effect = MaxRetriesExceededError monkeypatch.setattr(bulk_task, 'retry', retry) return retry def test_bulk_test_task(): """Test for testing that test task works well""" id = 'test_id' result = bulk.test_celery_task.apply(args=(id,)).get() assert result == '{0}: test passed.'.format(id) def test_lock_lock_isrc_exception( monkeypatch, sentry_captureException, client_error_fixture, feature_engine): """Test lock_isrc task for case when locks.lock_isrc raises exception """ lock_isrc = Mock() lock_isrc.side_effect = client_error_fixture monkeypatch.setattr(locks, 'lock_isrc', lock_isrc) mock_retry(monkeypatch, bulk.lock_isrc) isrc = 'US1' isrc_record = {} active_record = {} to_remove = ['QA', 'CA'] to_unlock = ['US', 'AF'] to_lock = ['DF'] correlation_id = str(uuid.uuid1()) user = 1 received_isrc, task_response = bulk.lock_isrc.delay( isrc, isrc_record, active_record, to_remove, to_unlock, to_lock, correlation_id, user).get() assert received_isrc == isrc assert isinstance(task_response, response.Response) assert task_response.status == 400 assert task_response.errors['code'] == error.LOCK_ISRC_ERROR_CODE assert task_response.errors['message'] == error.LOCK_ISRC_ERROR_MESSAGE lock_isrc.assert_called_once_with( isrc, isrc_record, active_record, to_remove, to_unlock, to_lock, correlation_id, user, None ) assert sentry_captureException.call_count == 1 def test_lock_lock_isrc(monkeypatch, feature_engine): """Test lock_isrc task """ lock_isrc = Mock(return_value=True) monkeypatch.setattr(locks, 'lock_isrc', lock_isrc) mock_retry(monkeypatch, bulk.lock_isrc) isrc = 'US1' isrc_record = {} active_record = {} to_remove = ['QA', 'CA'] to_unlock = ['US', 'AF'] to_lock = ['DF'] correlation_id = str(uuid.uuid1()) user = 1 result = bulk.lock_isrc.delay( isrc, isrc_record, active_record, to_remove, to_unlock, to_lock, correlation_id, user).get() assert result == (isrc, True) lock_isrc.assert_called_once_with( isrc, isrc_record, active_record, to_remove, to_unlock, to_lock, correlation_id, user, None ) def test_bulk_lock_isrcs_no_need_to_lock(monkeypatch): """Test bulk_lock_isrcs task for case where all territories are already locked and there is no need to make any changes in active_table """ monkeypatch.setattr( ownership, 'get_existing_isrcs_in_active_table', Mock(return_value=[])) monkeypatch.setattr( locks, '_check_territories', Mock(return_value=[[], [], []])) monkeypatch.setattr( locks, '_lock_create_active_table_records', Mock(return_value=[])) save_bulk_lock_report = Mock() monkeypatch.setattr( bulk.save_bulk_lock_report, 'delay', save_bulk_lock_report) correlation_id = str(uuid.uuid1()) user = '5' lock_reason = 'reason' isrcs = ['QA1', 'QA2'] territories_to_lock = ['QA', 'CA'] task_id = 1 result = [('QA1', True), ('QA2', True)] bulk.bulk_lock_isrcs( lock_reason, isrcs, territories_to_lock, correlation_id, user, task_id) save_bulk_lock_report.assert_called_with( result, task_id, territories_to_lock, lock_reason, 0 ) def test_bulk_lock_isrcs(monkeypatch): """Test bulk_lock_isrcs""" isrc_record = { 'isrc': 'QA1', 'locked_territories': {'CA': {'reason': 'foo'}}, 'territories': {} } monkeypatch.setattr( ownership, 'get_existing_isrcs_in_active_table', Mock(return_value=[isrc_record]) ) monkeypatch.setattr(bulk.save_bulk_lock_report, 's', Mock()) monkeypatch.setattr( bulk, 'chord', Mock()) lock_isrc = Mock() monkeypatch.setattr(bulk.lock_isrc, 's', lock_isrc) correlation_id = str(uuid.uuid1()) user = '5' lock_reason = 'reason' isrcs = ['QA1'] territories_to_lock = ['QA', 'CA'] task_id = 1 bulk.bulk_lock_isrcs( lock_reason, isrcs, territories_to_lock, correlation_id, user, task_id) lock_isrc.assert_called_with( 'QA1', isrc_record, ANY, set(), {'CA'}, {'CA', 'QA'}, correlation_id, user, territories_to_lock ) def test_save_bulk_lock_report_error( mocker, monkeypatch, sentry_captureMessage, default_task_data): """Expect error from update_task.""" # TODO refactor to use database instead of mocks? task_data = default_task_data get_task_response = response.Response(bulk_tasks.TaskStatus(**task_data)) mocker.patch.object( bulk_tasks, 'get_task', return_value=get_task_response) update_task_response = response.create_error_response( code='error', message='Error', status=500) mocker.patch.object( bulk_tasks, 'update_task', return_value=update_task_response) mocker.patch.object(bulk, 'generate_reports') retry = mock_retry(monkeypatch, bulk.save_bulk_lock_report) bulk.save_bulk_lock_report.delay('', '', '', '', '').get() assert retry.call_count == 1 assert sentry_captureMessage.call_count == 1 def test_unlock_lock_isrc_exception( monkeypatch, sentry_captureException, client_error_fixture): """Test unlock_isrc task for case when locks.lock_isrc raises exception """ unlock_isrc = Mock() unlock_isrc.side_effect = client_error_fixture monkeypatch.setattr(locks, 'unlock_isrc', unlock_isrc) mock_retry(monkeypatch, bulk.unlock_isrc) isrc = 'US1' territories = ['US', 'AF'] correlation_id = str(uuid.uuid1()) user = 1 result = bulk.unlock_isrc.delay( isrc, territories, correlation_id, user).get() assert result == (isrc, False) unlock_isrc.assert_called_once_with( isrc, territories, correlation_id, user) assert sentry_captureException.call_count == 1 def test_unlock_lock_isrc(monkeypatch): """Test unlock_isrc task """ unlock_isrc = Mock(return_value=True) monkeypatch.setattr(locks, 'unlock_isrc', unlock_isrc) mock_retry(monkeypatch, bulk.lock_isrc) isrc = 'US1' territories = ['US', 'AF'] correlation_id = str(uuid.uuid1()) user = 1 result = bulk.unlock_isrc.delay( isrc, territories, correlation_id, user).get() assert result == (isrc, True) unlock_isrc.assert_called_once_with( isrc, territories, correlation_id, user) def test_save_bulk_unlock_report_partial_fail(monkeypatch): """Test test_save_bulk_unlock_report task for partial fail case """ update_task = Mock() monkeypatch.setattr(bulk_tasks, 'update_task', update_task) unlock_result = [(1, False), (2, True)] task_id = 1 bulk.update_bulk_unlock_status.apply( args=(unlock_result, task_id, 2)).get() update_task.assert_called_with( task_id, bulk_tasks_const.PARTIAL_FAIL_STATUS, None ) def test_save_bulk_unlock_report_fail(monkeypatch): """Test test_save_bulk_unlock_report task for fail case """ update_task = Mock() monkeypatch.setattr(bulk_tasks, 'update_task', update_task) unlock_result = [(1, False)] task_id = 2 bulk.update_bulk_unlock_status.apply( args=(unlock_result, task_id, 1)).get() update_task.assert_called_with( task_id, bulk_tasks_const.FAILED_STATUS, None ) def test_save_bulk_unlock_report_error(monkeypatch, sentry_captureMessage): """Test bulk_unlock_report task for case when bulk_tasks.update_task returns error """ update_task = Mock(return_value=response.create_error_response( code='error', message='Error', status=500 )) monkeypatch.setattr(bulk_tasks, 'update_task', update_task) retry = mock_retry(monkeypatch, bulk.update_bulk_unlock_status) bulk.update_bulk_unlock_status.delay('', '', '').get() assert retry.call_count == 1 assert sentry_captureMessage.call_count == 1 @patch( 'masters_registry.models.ows_territories.convert_territories', new=patches.ows_territories_convert_territories_patch) @patch( 'masters_registry.models.ownership.get_upc_by_tuid', new=patches.ownership_get_upc_by_tuid_patch) @patch( 'masters_registry.models.ows_carveouts.get_dms_carveout_for_upc', new=patches.ows_carveouts_get_dms_carveout_for_upc_not_carved_out_patch) def test_update_yt_for_deleted_upc(monkeypatch): """Test update_yt_for_deleted_upc task""" correlation_id = '1234567890' isrc = 'US94C0519904' ownership_info = { 'isrc': isrc, 'locked_territories': { 'AD': {'reason': 'lock reason'} }, 'territories': { 'QA': {'tuid': '123'}, 'AF': {'tuid': '123'}, 'CF': {'tuid': '789'}, } } removed_territories = ['QA', 'AF'] update_result = { field_const.OWNERSHIP_INFO: ownership_info, field_const.REMOVED_TERRITORIES: removed_territories, field_const.FAILED: False } expected_result = (True, '') send_message = Mock() monkeypatch.setattr( yt_ownership, 'send_message', send_message) result = bulk.update_yt_for_deleted_upc( update_result, isrc, correlation_id) send_message.assert_called_with(isrc, ['CF'], correlation_id) assert result == expected_result def test_update_yt_for_deleted_upc_empty_updated_claimed_territories( monkeypatch): """Test update_yt_for_deleted_upc task for case when all claimed territories are removed """ correlation_id = '1234567890' isrc = 'US94C0519904' ownership_info = { 'isrc': isrc, 'locked_territories': { 'AD': {'reason': 'lock reason'} }, 'territories': { 'QA': {'tuid': '123'}, 'AF': {'tuid': '123'}, } } removed_territories = ['QA', 'AF'] update_result = { field_const.OWNERSHIP_INFO: ownership_info, field_const.REMOVED_TERRITORIES: removed_territories, field_const.FAILED: False } expected_result = (True, '') send_message = Mock() monkeypatch.setattr( yt_ownership, 'send_message', send_message) result = bulk.update_yt_for_deleted_upc( update_result, isrc, correlation_id) send_message.assert_called_with(isrc, set(), correlation_id) assert result == expected_result def test_update_yt_for_deleted_upc_no_store_carveout( monkeypatch): """Test update_yt_for_deleted_upc task for case when ownership_check_dms_carveout returns error response """ correlation_id = '1234567890' isrc = 'US94C0519904' ownership_info = { 'isrc': isrc, 'locked_territories': { 'AD': {'reason': 'lock reason'} }, 'territories': { 'QA': {'tuid': '123'}, 'AF': {'tuid': '123'}, 'CF': {'tuid': '789'}, } } removed_territories = ['QA', 'AF'] update_result = { field_const.OWNERSHIP_INFO: ownership_info, field_const.REMOVED_TERRITORIES: removed_territories, field_const.FAILED: False } expected_result = (False, error.FAILED_TO_GET_DMS_CARVEOUTS) ownership_check_dms_carveout = Mock( return_value=response.create_error_response( code='error', message='Error' ) ) monkeypatch.setattr( dms_carveout, 'ownership_check_dms_carveout', ownership_check_dms_carveout) result = bulk.update_yt_for_deleted_upc( update_result, isrc, correlation_id) assert result == expected_result def test_update_yt_for_deleted_upc_upc_carved_out(monkeypatch): """Test update_yt_for_deleted_upc task for case when UPC is carved out from YouTube """ correlation_id = '1234567890' isrc = 'US94C0519904' ownership_info = { 'isrc': isrc, 'locked_territories': { 'AD': {'reason': 'lock reason'} }, 'territories': { 'QA': {'tuid': '123'}, 'AF': {'tuid': '123'}, 'CF': {'tuid': '789'}, } } removed_territories = ['QA', 'AF'] update_result = { field_const.OWNERSHIP_INFO: ownership_info, field_const.REMOVED_TERRITORIES: removed_territories, field_const.FAILED: False } expected_result = (False, error.UPC_YOUTUBE_CARVED_OUT_MESSAGE) ownership_check_dms_carveout = Mock(return_value=response.Response( message={} )) monkeypatch.setattr( dms_carveout, 'ownership_check_dms_carveout', ownership_check_dms_carveout) result = bulk.update_yt_for_deleted_upc( update_result, isrc, correlation_id) assert result == expected_result def test_save_bulk_import_report_done( generate_reports, monkeypatch, valid_correlation_id, feature_engine): """Test save_bulk_import_report task for success case""" update_task = Mock() monkeypatch.setattr(bulk_tasks, 'update_task', update_task) import_results = [ { field_const.UPC: '1', field_const.ISRC: 'US1', field_const.SUCCESS: True, field_const.ERROR_MESSAGE: '', field_const.LOCKED_TERRITORIES: ['US'] }, { field_const.UPC: '2', field_const.ISRC: 'CA1', field_const.SUCCESS: True, field_const.ERROR_MESSAGE: '', field_const.LOCKED_TERRITORIES: ['CA'] }, { field_const.UPC: '2', field_const.ISRC: 'CA2', field_const.SUCCESS: True, field_const.ERROR_MESSAGE: '', field_const.LOCKED_TERRITORIES: ['CA', 'US'] } ] import_upc_report = [ { field_const.UPC: '1', field_const.ERROR_COUNT: 0, field_const.WARNING_COUNT: 0, field_const.SUCCESS_COUNT: 1, field_const.STATUS_REPORT: [ { field_const.UPC: '1', field_const.ISRC: 'US1', field_const.SUCCESS: True, field_const.LOCKED_TERRITORIES: ['US'], field_const.ERROR_MESSAGE: '', } ], }, { field_const.UPC: '2', field_const.ERROR_COUNT: 0, field_const.WARNING_COUNT: 0, field_const.SUCCESS_COUNT: 2, field_const.STATUS_REPORT: [ { field_const.UPC: '2', field_const.ISRC: 'CA1', field_const.SUCCESS: True, field_const.LOCKED_TERRITORIES: ['CA'], field_const.ERROR_MESSAGE: '' }, { field_const.UPC: '2', field_const.ISRC: 'CA2', field_const.SUCCESS: True, field_const.LOCKED_TERRITORIES: ['CA', 'US'], field_const.ERROR_MESSAGE: '', } ] } ] task_id = 1 expected_status = bulk_tasks_const.DONE_STATUS user = '1234' bulk.save_bulk_import_report.apply( (import_results, task_id, 2, user, valid_correlation_id)).get() update_task.assert_called_once_with( task_id, expected_status, import_upc_report ) generate_reports.assert_called_once_with( task_id, [field_const.ISRC, field_const.UPC], valid_correlation_id ) def test_save_bulk_import_report_partial_fail( generate_reports, monkeypatch, valid_correlation_id, feature_engine): """Test save_bulk_import_report task for partial fail case""" update_task = Mock() monkeypatch.setattr(bulk_tasks, 'update_task', update_task) import_results = [ { field_const.UPC: '1', field_const.ISRC: 'US1', field_const.SUCCESS: False, field_const.ERROR_MESSAGE: 'Error', field_const.LOCKED_TERRITORIES: ['US'] }, { field_const.UPC: '1', field_const.ISRC: 'CA1', field_const.SUCCESS: True, field_const.ERROR_MESSAGE: '', field_const.LOCKED_TERRITORIES: ['CA'] } ] import_upc_report = [ { field_const.UPC: '1', field_const.ERROR_COUNT: 1, field_const.SUCCESS_COUNT: 1, field_const.WARNING_COUNT: 0, field_const.STATUS_REPORT: [ { field_const.UPC: '1', field_const.ISRC: 'US1', field_const.SUCCESS: False, field_const.LOCKED_TERRITORIES: ['US'], field_const.ERROR_MESSAGE: 'Error' }, { field_const.UPC: '1', field_const.ISRC: 'CA1', field_const.SUCCESS: True, field_const.LOCKED_TERRITORIES: ['CA'], field_const.ERROR_MESSAGE: '' }, ] } ] task_id = 1 expected_status = bulk_tasks_const.PARTIAL_FAIL_STATUS user = '1234' bulk.save_bulk_import_report.apply( (import_results, task_id, 2, user, valid_correlation_id)).get() update_task.assert_called_once_with( task_id, expected_status, import_upc_report ) generate_reports.assert_called_once_with( task_id, [field_const.ISRC, field_const.UPC], valid_correlation_id ) def test_save_bulk_import_report_fail( generate_reports, monkeypatch, valid_correlation_id, feature_engine): """Test save_bulk_import_report task for fail case""" update_task = Mock() monkeypatch.setattr(bulk_tasks, 'update_task', update_task) import_results = [{ field_const.UPC: '1', field_const.ISRC: 'US1', field_const.SUCCESS: False, field_const.ERROR_MESSAGE: 'Error', field_const.LOCKED_TERRITORIES: ['US'] }] import_upc_report = [ { field_const.UPC: '1', field_const.ERROR_COUNT: 1, field_const.SUCCESS_COUNT: 0, field_const.WARNING_COUNT: 0, field_const.STATUS_REPORT: [ { field_const.UPC: '1', field_const.ISRC: 'US1', field_const.SUCCESS: False, field_const.LOCKED_TERRITORIES: ['US'], field_const.ERROR_MESSAGE: 'Error' }, ] } ] task_id = 1 expected_status = bulk_tasks_const.FAILED_STATUS user = '1234' bulk.save_bulk_import_report.apply( (import_results, task_id, 1, user, valid_correlation_id)).get() update_task.assert_called_once_with( task_id, expected_status, import_upc_report ) generate_reports.assert_called_once_with( task_id, [field_const.ISRC, field_const.UPC], valid_correlation_id ) def test_save_bulk_import_report_error( monkeypatch, sentry_captureMessage, valid_correlation_id, feature_engine): """Test save_bulk_import_report task for case when bulk_tasks.update_task returns error """ import_results = [{ field_const.UPC: '1', field_const.ISRC: 'US1', field_const.SUCCESS: False, field_const.ERROR_MESSAGE: 'Error', field_const.LOCKED_TERRITORIES: ['US'] }] user = '1234' update_task = Mock(return_value=response.create_error_response( code='error', message='Error', status=500 )) monkeypatch.setattr(bulk_tasks, 'update_task', update_task) retry = mock_retry(monkeypatch, bulk.save_bulk_import_report) bulk.save_bulk_import_report.delay( import_results, 1, 1, user, valid_correlation_id).get() assert retry.call_count == 1 assert sentry_captureMessage.call_count == 1 @patch('masters_registry.models.ownership.get_isrcs', new=patches.ownership_get_isrcs_bulk_import_patch) @patch('masters_registry.models.ownership.get_upc_by_tuid', new=patches.ownership_get_upc_by_tuid_patch) @patch('masters_registry.models.yt_ownership.send_message', new=Mock()) @patch('masters_registry.models.ows_carveouts.get_dms_carveout_for_upc', new=patches.ows_carveouts_get_dms_carveout_for_upc_not_carved_out_patch) @patch('masters_registry.models.ownership.get_tracks', new=patches.ownership_get_tracks_patch) @patch('masters_registry.logic.masters_registry._get_carved_in_territories', new=Mock(return_value=response.Response(message=['US', 'CA']))) @patch('masters_registry.models.ows_territories.convert_territories', new=patches.ows_territories_convert_territories_patch) @patch('masters_registry.logic.bulk_tasks.upload_report', new=Mock()) @mock_aws def test_bulk_import_upcs( setup_masters_active_table, seed_masters_active_table, setup_new_audit_table, setup_bulk_status_db, feature_engine): """Test bulk_import_upcs task""" from moto.core import patch_client, patch_resource patch_client(outside_dynamodb_client) patch_resource(dynamodb_resource) ownership_info = bulk_import_fixtures.ownership_info_conflicts_enabled() expected_active_table_results = ( bulk_import_fixtures.active_table_conflicts_enabled()) expected_audit_table_results = ( bulk_import_fixtures.audit_table_conflicts_enabled()) expected_task_status = bulk_tasks_const.WARNING_STATUS expected_task_result = bulk_import_fixtures.task_result_conflicts_enabled() setup_masters_active_table(dynamodb_resource.meta.client) setup_new_audit_table(dynamodb_resource.meta.client) for _ownership in ownership_info: seed_masters_active_table(dynamodb_resource.meta.client, item=_ownership) upcs = ['886788771130', '767767806329'] correlation_id = '123' user = '1234' ignore_keys = ['timestamp', 'updated_timestamp'] task_id = bulk_tasks.create_task( correlation_id, user, bulk_tasks_const.BULK_RESOLVE_INTERNAL_CONFLICTS, len(upcs), upcs) bulk.bulk_import_upcs.apply((upcs, task_id, correlation_id, user)) send_message_calls = [ call('USJTQ1109709', set(), correlation_id), call('USA371690671', ['CA', 'QA', 'US'], correlation_id), ] assert yt_ownership.send_message.call_args_list == send_message_calls yt_ownership.send_message.reset_mock() for isrc, expected_results in expected_audit_table_results.items(): audit_records = ownership.get_ownership_audit(isrc).message for i, record in enumerate(audit_records): assert equal_dicts(record, expected_results[i], ignore_keys) for isrc, expected_result in expected_active_table_results.items(): assert equal_dicts( ownership.get_ownership(isrc).message, expected_result, ignore_keys) task_result = bulk_tasks.get_task(task_id).message assert task_result.status == expected_task_status assert task_result.result == expected_task_result def test_update_isrc_for_upc_exception( monkeypatch, sentry_captureException, client_error_fixture): """Test update_isrc_for_upc task for case when _process_isrc_info raises exception """ _process_isrc_info = Mock() _process_isrc_info.side_effect = client_error_fixture monkeypatch.setattr( masters_registry, '_process_isrc_info', _process_isrc_info) retry = mock_retry(monkeypatch, bulk.update_isrc_for_upc) upc = 1 isrc_info = {field_const.ISRC: 'US1'} carve_in_territories = response.Response(message='') bulk.update_isrc_for_upc.delay( upc, isrc_info, carve_in_territories, '', '').get() assert retry.call_count == 1 assert sentry_captureException.call_count == 1 def test_update_isrc_for_upc_key_error_exception( monkeypatch, sentry_captureException, client_error_fixture): """Test update_isrc_for_upc task for case when _process_isrc_info raises key error exception """ _process_isrc_info = Mock() _process_isrc_info.side_effect = KeyError() monkeypatch.setattr( masters_registry, '_process_isrc_info', _process_isrc_info) retry = mock_retry(monkeypatch, bulk.update_isrc_for_upc) upc = 1 isrc_info = {field_const.ISRC: 'US1'} carve_in_territories = response.Response(message='') bulk.update_isrc_for_upc.delay( upc, isrc_info, carve_in_territories, '', '').get() assert retry.call_count == 0 assert sentry_captureException.call_count == 1 @mock_aws def test_update_isrc_conflict_resolved( setup_masters_active_table, seed_masters_active_table, setup_new_audit_table): """Test update_isrc_for_upc task for case when an internal conflict was resolved.""" from moto.core import patch_client, patch_resource patch_client(outside_dynamodb_client) patch_resource(dynamodb_resource) patch_client(sqs_client) patch_resource(sqs_resource) sqs_client.create_queue(QueueName=config.QUEUE_NAME_YT_OWNERSHIP) queue_yt_ownership = sqs_resource.get_queue_by_name( QueueName=config.QUEUE_NAME_YT_OWNERSHIP) flexmock(yt_ownership).should_receive( 'sqs.queue_yt_ownership').and_return(queue_yt_ownership) upc = '987654321' tuid = 321 vendor_id = 7 correlation_id = str(uuid.uuid1()) isrc = '131131131' user = '123' ownership_info = { field_const.ISRC: isrc, field_const.TUID: tuid, field_const.LOCKED_TERRITORIES: {}, field_const.TERRITORIES: { 'QA': [{field_const.TUID: tuid}, {field_const.TUID: '123'}], 'CA': {field_const.TUID: tuid}, } } setup_masters_active_table(dynamodb_resource.meta.client) setup_new_audit_table(dynamodb_resource.meta.client) seed_masters_active_table(dynamodb_resource.meta.client, item=ownership_info) result = bulk.update_isrc_for_deleted_upc.delay( upc, isrc, tuid, vendor_id, user, correlation_id).get() expected_result = { field_const.UPC: upc, field_const.ISRC: isrc, field_const.TUID: tuid, field_const.VENDOR_ID: vendor_id, field_const.SUCCESS: True, field_const.WARNING: True, field_const.RESOLVED_CONFLICT: True, field_const.REMOVED_TERRITORIES: ['CA', 'QA'], field_const.ERROR_MESSAGE: '', } assert result == expected_result def test_update_isrc_for_deleted_upc_exception( monkeypatch, sentry_captureException, client_error_fixture): """Test update_isrc_for_deleted_upc task for case when ownership.get_ownership raise exception """ get_ownership = Mock() get_ownership.side_effect = client_error_fixture monkeypatch.setattr(ownership, 'get_ownership', get_ownership) retry = mock_retry(monkeypatch, bulk.update_isrc_for_deleted_upc) bulk.update_isrc_for_deleted_upc.delay('', '', '', '', '', '').get() assert retry.call_count == 1 assert sentry_captureException.call_count == 1 def test_update_isrc_for_upc_no_isrc(valid_correlation_id): """Test update_isrc_for_upc task for case when isrc_info doesn't have ISRC """ upc = 1 isrc_info = {field_const.ISRC: None} carve_in_territories = response.Response(message='') failed_response = { field_const.UPC: upc, field_const.ISRC: None, field_const.TUID: None, field_const.VENDOR_ID: None, field_const.SUCCESS: False, field_const.ERROR_MESSAGE: error.NO_ISRC_ERROR } result = bulk.update_isrc_for_upc.delay( upc, isrc_info, carve_in_territories, '', '').get() assert result == failed_response def test_generate_reports(generate_report_logic, valid_correlation_id): """Test generate_reports function""" task_id = 69 bulk.generate_reports( task_id, [field_const.ISRC, field_const.UPC], valid_correlation_id) generate_report_logic.assert_has_calls(( call( task_id, valid_correlation_id, import_report_type=field_const.ISRC), call( task_id, valid_correlation_id, import_report_type=field_const.UPC), )) def test_generate_reports_exception( sentry_captureException, client_error_fixture, monkeypatch, valid_correlation_id): """Test generate_reports function for a case when generate_report raises an exception """ generate_report = Mock() generate_report.side_effect = client_error_fixture monkeypatch.setattr(bulk_tasks_logic, 'generate_report', generate_report) bulk.generate_reports( 12, [field_const.ISRC, field_const.UPC], valid_correlation_id) assert sentry_captureException.call_count == 1 @patch('masters_registry.models.bulk_tasks.update_task', new=Mock()) @pytest.mark.parametrize('successful,failed,expected_status', [ ({'US1': {}, 'US2': {}}, {'CA': {}}, bulk_tasks_const.PARTIAL_FAIL_STATUS), ({}, {'CA1': {}}, bulk_tasks_const.FAILED_STATUS), ({'US1': {}}, {}, bulk_tasks_const.DONE_STATUS), ]) def test_save_resolve_conflict_result_partial_fail( successful, failed, expected_status): """Test save_resolve_conflict_result task """ territories = ['QA', 'CA'] account_id = '123' task_id = 1 task_result = { field_const.TERRITORIES: territories, field_const.SUCCESSFUL_ISRCS: successful, field_const.FAILED_ISRCS: failed, field_const.ACCOUNT_ID: account_id } bulk.save_resolve_conflict_result.apply( args=(territories, successful, failed, task_id, account_id)).get() bulk_tasks.update_task.assert_called_with( task_id, expected_status, task_result ) @patch('masters_registry.models.bulk_tasks.update_task', new=Mock( return_value=response.create_error_response( code='error', message='Error', status=500))) def test_save_resolve_conflict_result_error( monkeypatch, sentry_captureMessage): """Test save_resolve_conflict_result task for case when bulk_tasks.update_task returns error """ retry = mock_retry(monkeypatch, bulk.save_resolve_conflict_result) bulk.save_resolve_conflict_result.delay([], {}, {}, '', '').get() assert retry.call_count == 1 assert sentry_captureMessage.call_count == 1 NOT_WORKING_ISRC = 'NOT_WORKING' def mock_remove_ownership( isrc, tuid, territories_to_remove, existing_territories, correlation_id, user, daemon=False, source=''): if isrc == NOT_WORKING_ISRC: raise client_error() return { field_const.UPDATED_CLAIMED_TERRITORIES: territories_to_remove, field_const.RESOLVED_CONFLICT: {} } @patch('time.sleep', new=Mock()) @patch( 'masters_registry.models.ownership.get_ownership', new=patches.ownership_get_ownership_with_conflicts_patch) @patch( 'masters_registry.logic.ownership.remove_ownership', new=Mock(side_effect=mock_remove_ownership)) @patch('masters_registry.tasks.bulk.save_resolve_conflict_result.delay', new=Mock()) @patch( 'masters_registry.models.yt_ownership.send_message', new=Mock()) def test_bulk_resolve_conflicts_remove_ownership_exception( sentry_captureException): """Test bulk_resolve_conflicts for a case when remove_ownership raises an exception. """ isrc_tuid_map = { 'CA1': [111, 222], NOT_WORKING_ISRC: [111] } territories = { 'CA1': ['CA', 'PT', 'QA', 'MX'], NOT_WORKING_ISRC: ['CA', 'PT', 'QA', 'MX']} orchard_user_id = '7788' task_id = 100 correlation_id = '123' account_id = '987' successful = {'CA1': {222: {'PT'}, 111: {'QA'}}} failed = {NOT_WORKING_ISRC: [111]} existing_territories = ( patches.ownership_get_ownership_with_conflicts_patch('') .message[field_const.TERRITORIES]) bulk.bulk_resolve_conflicts.delay( isrc_tuid_map, territories, orchard_user_id, task_id, correlation_id, account_id) bulk.save_resolve_conflict_result.delay.assert_called_with( {'CA', 'PT', 'QA', 'MX'}, successful, failed, task_id, account_id) ownership_logic.remove_ownership.assert_any_call( 'CA1', 111, ['QA'], existing_territories, correlation_id, orchard_user_id, daemon=True, source=field_const.BULK_RESOLVE_CONFLICTS ) ownership_logic.remove_ownership.assert_any_call( 'CA1', 222, ['PT'], existing_territories, correlation_id, orchard_user_id, daemon=True, source=field_const.BULK_RESOLVE_CONFLICTS ) assert sentry_captureException.call_count == 3 @patch('masters_registry.models.ownership.get_tracks', new=patches.ownership_get_tracks_patch) @mock_aws def test_bulk_resolve_conflicts( setup_bulk_status_db, setup_masters_active_table, seed_masters_active_table, setup_new_audit_table, feature_engine): """Functional test for bulk_resolve_conflicts""" from moto.core import patch_client, patch_resource patch_client(outside_dynamodb_client) patch_resource(dynamodb_resource) patch_client(outside_s3_client) patch_resource(s3_resource) bucket_name = config.REPORTS_BUCKET_NAME outside_s3_client.create_bucket(Bucket=bucket_name) ownership_info = { field_const.ISRC: 'CA1', field_const.LOCKED_TERRITORIES: {}, field_const.TERRITORIES: { 'CA': {field_const.TUID: 12345}, 'MX': [{field_const.TUID: 11}], 'QA': [ {field_const.TUID: 11}, {field_const.TUID: 12345}, ], 'PT': [ {field_const.TUID: 222}, {field_const.TUID: 123456}, ] } } setup_masters_active_table(dynamodb_resource.meta.client) seed_masters_active_table(dynamodb_resource.meta.client, item=ownership_info) setup_new_audit_table(dynamodb_resource.meta.client) test_isrc = 'CA1' isrc_tuid_map = {test_isrc: [11, 222]} territories = {test_isrc: ['CA', 'PT', 'QA', 'MX']} orchard_user_id = '7788' correlation_id = '123' account_id = '987' task_id = bulk_tasks.create_task( correlation_id, orchard_user_id, bulk_tasks_const.BULK_RESOLVE_INTERNAL_CONFLICTS, 1, [test_isrc] ) bulk.bulk_resolve_conflicts.delay( isrc_tuid_map, territories, orchard_user_id, task_id, correlation_id, account_id) saved_item = ownership.get_ownership(test_isrc).message expected_saved_item = { 'isrc': test_isrc, 'locked_territories': {}, 'territories': { 'MX': [{'tuid': Decimal('11')}], 'QA': [{'tuid': Decimal('12345')}], 'CA': {'tuid': Decimal('12345')}, 'PT': [{'tuid': Decimal('123456')}] }, } saved_item.pop('updated_timestamp') assert saved_item == expected_saved_item expected_audit_records = [ { 'correlation_id': correlation_id, 'user': orchard_user_id, 'isrc': test_isrc, 'opcode': opcode_const.REMOVE, 'territories': {'QA': Decimal('11')} }, { 'correlation_id': correlation_id, 'user': orchard_user_id, 'isrc': test_isrc, 'source': field_const.BULK_RESOLVE_CONFLICTS, 'opcode': opcode_const.CONFLICT_RESOLVED, 'conflict': { 'resolved': Decimal('1'), 'status': opcode_const.CONFLICT_RESOLVED, 'conflicting_tuid': [Decimal('123456'), Decimal('222')] }, 'territories': {'QA': Decimal('11')} }, { 'correlation_id': correlation_id, 'user': orchard_user_id, 'isrc': test_isrc, 'opcode': opcode_const.REMOVE, 'territories': {'PT': Decimal('222')} }, { 'correlation_id': correlation_id, 'user': orchard_user_id, 'isrc': test_isrc, 'source': field_const.BULK_RESOLVE_CONFLICTS, 'opcode': opcode_const.CONFLICT_RESOLVED, 'conflict': { 'resolved': Decimal('1'), 'status': opcode_const.CONFLICT_RESOLVED, 'conflicting_tuid': [] }, 'territories': {'PT': Decimal('222')} } ] audit_records = ownership.get_ownership_audit(test_isrc).message for record in audit_records: record.pop('timestamp') assert audit_records == expected_audit_records expected_task_result = { field_const.FAILED_ISRCS: {}, field_const.SUCCESSFUL_ISRCS: {'CA1': {11: {'QA'}, 222: {'PT'}}}, field_const.TERRITORIES: {'PT', 'QA', 'MX', 'CA'}, field_const.ACCOUNT_ID: account_id } task_results = bulk_tasks.get_task(task_id).message assert task_results.status == bulk_tasks_const.DONE_STATUS assert task_results.result == expected_task_result date = str(task_results.create_datetime.date()) file_key = '{0}/{1}_Bulk_Resolve_Conflicts_ISRC_{2}.csv'.format( config.REPORTS_FILE_KEY_PREFIX, date, task_id) assert s3.check_if_file_exists(file_key) @patch('masters_registry.models.ownership.update_active_record', new=Mock()) @patch('masters_registry.logic.locks.lock_update_audit_table', new=Mock()) @patch( 'masters_registry.logic.locks.lock_update_yt_claimed_territories', new=Mock(return_value=[])) def test_lock_isrc_internal_conflict_oto_response_success(): """Expect to get correct task.""" isrc = 'QA123' isrc_record = { 'isrc': 'QA123', 'territories': {'US': {'tuid': 123}}, 'locked_territories': {'PL': {'reason': 'Owned by Warner'}} } reason = 'Lock reason' to_remove = [] to_unlock = [] to_lock = ['US'] correlation_id = '123' user = '1' active_record = { 'territories': { 'to_lock': to_lock, 'to_remove': to_remove, 'to_unlock': to_unlock, }, 'lock_reason': reason, 'isrc': 'QA123', } returned_isrc, task_response = bulk.lock_isrc.delay( isrc, isrc_record, active_record, to_remove, to_unlock, to_lock, correlation_id, user, to_lock).get() assert returned_isrc == isrc assert isinstance(task_response, response.Response) assert task_response.status == 200 ownership.update_active_record.assert_called_once_with( 'QA123', {'to_unlock': [], 'to_lock': ['US'], 'to_remove': ['US']}, lock_reason='Lock reason') locks.lock_update_audit_table.assert_called_once_with( ({'QA123': {'US'}}, 'REMOVE'), ({'QA123': set()}, 'UNLOCK'), ({'QA123': {'US'}}, 'LOCK'), '123', '1', 'Lock reason') locks.lock_update_yt_claimed_territories.assert_called_once_with( [{'isrc': 'QA123', 'locked_territories': {'PL': {'reason': 'Owned by Warner'}}, 'territories': {'US': {'tuid': 123}}}], {'QA123': {'US'}}, '123') @patch('masters_registry.models.ownership.update_active_record', new=Mock()) @patch('masters_registry.logic.locks.lock_update_audit_table', new=Mock()) @patch( 'masters_registry.logic.locks.lock_update_yt_claimed_territories', new=Mock(return_value=[])) def test_lock_isrc_internal_conflict_oto_response_failure(): """Should not lock territory that has internal conflict.""" isrc = 'QA123' isrc_record = { 'isrc': 'QA123', 'territories': {'US': [{'tuid': 123}, {'tuid': 456}]}, 'locked_territories': {'PL': {'reason': 'Owned by Warner'}} } reason = 'Lock reason' to_remove = [] to_unlock = [] to_lock = ['US'] correlation_id = '123' user = '1' active_record = { 'territories': { 'to_lock': to_lock, 'to_remove': to_remove, 'to_unlock': to_unlock, }, 'lock_reason': reason, 'isrc': 'QA123', } returned_isrc, task_response = bulk.lock_isrc.delay( isrc, isrc_record, active_record, to_remove, to_unlock, to_lock, correlation_id, user, to_lock).get() assert returned_isrc == isrc assert isinstance(task_response, response.Response) assert task_response.status == 400 expected_message = {'territories': [], 'failed_territories': ['US']} assert task_response.errors['message'] == expected_message ownership.update_active_record.assert_called_once_with( 'QA123', {'to_remove': [], 'to_unlock': [], 'to_lock': []}, lock_reason='Lock reason') locks.lock_update_audit_table.assert_called_once_with( ({'QA123': set()}, 'REMOVE'), ({'QA123': set()}, 'UNLOCK'), ({'QA123': set()}, 'LOCK'), '123', '1', 'Lock reason') locks.lock_update_yt_claimed_territories.assert_called_once_with( [{'isrc': 'QA123', 'locked_territories': {'PL': {'reason': 'Owned by Warner'}}, 'territories': {'US': [{'tuid': 123}, {'tuid': 456}]}}], {'QA123': set()}, '123') def test_make_lock_result_internal_conflict_success(): """Should make report data without internal conflict.""" lock_isrc_results = [('QA123', response.Response())] lock_result = bulk.make_lock_result_internal_conflict( lock_isrc_results, territories=['AF'], reason='test') expected_lock_result = { 'isrcs': { 'QA123': {'territories': ['AF'], 'failed_territories': []}}, 'reason': 'test', 'failed_isrcs': False, 'successful_isrcs': True } assert lock_result == expected_lock_result def test_make_lock_result_internal_conflict_error(): """Should make report data with internal conflict.""" failed_isrcs = response.create_error_response( code='', message={ 'territories': ['AF'], 'failed_territories': ['AD'] }) lock_isrc_results = [('QA123', failed_isrcs)] lock_result = bulk.make_lock_result_internal_conflict( lock_isrc_results, territories=['AF', 'AD'], reason='test') expected_lock_result = { 'isrcs': { 'QA123': {'territories': ['AF'], 'failed_territories': ['AD']}}, 'reason': 'test', 'failed_isrcs': True, 'successful_isrcs': True } assert lock_result == expected_lock_result @patch('masters_registry.tasks.bulk.generate_reports', new=Mock()) def test_save_bulk_lock_report_no_internal_conflict( setup_bulk_status_db, seed_bulk_status_db, default_task_data, feature_engine): """Expect to save bulk lock report without internal conflict.""" seed_bulk_status_db([bulk_tasks.TaskStatus(**default_task_data)]) lock_isrc_result = [('QA123', response.Response())] bulk.save_bulk_lock_report.delay( lock_isrc_results=lock_isrc_result, task_id=1, territories=['AF'], reason='test', lock_tasks_amount=1) updated_task = bulk_tasks.get_task(1).message.as_dict() expected_task_result = { 'failed_isrcs': False, 'successful_isrcs': True, 'isrcs': {'QA123': {'territories': ['AF'], 'failed_territories': []}}, 'reason': 'test' } assert updated_task['result'] == expected_task_result assert updated_task['status'] == 'DONE' bulk.generate_reports.assert_called_once_with(1, ['isrc']) @patch('masters_registry.tasks.bulk.generate_reports', new=Mock()) def test_save_bulk_lock_report_with_internal_conflict( setup_bulk_status_db, seed_bulk_status_db, default_task_data, feature_engine): """Expect to save bulk lock report with internal conflict.""" seed_bulk_status_db([bulk_tasks.TaskStatus(**default_task_data)]) failed_isrcs = response.create_error_response( code='', message={ 'territories': ['AF'], 'failed_territories': ['AD'] }) lock_isrc_result = [('QA123', failed_isrcs)] bulk.save_bulk_lock_report.delay( lock_isrc_results=lock_isrc_result, task_id=1, territories=['AF', 'AD'], reason='test', lock_tasks_amount=1) updated_task = bulk_tasks.get_task(1).message.as_dict() expected_task_result = { 'reason': 'test', 'successful_isrcs': True, 'failed_isrcs': True, 'isrcs': { 'QA123': {'failed_territories': ['AD'], 'territories': ['AF']}} } assert updated_task['result'] == expected_task_result assert updated_task['status'] == 'PARTIAL_FAIL' bulk.generate_reports.assert_called_once_with(1, ['isrc']) @patch('masters_registry.tasks.bulk.generate_reports', new=Mock()) def test_save_bulk_lock_report_with_internal_conflict_no_success( setup_bulk_status_db, seed_bulk_status_db, default_task_data, feature_engine): """Expect to save bulk lock report with internal conflict. Should handle fact that there were no successfully locked ISRCs. """ seed_bulk_status_db([bulk_tasks.TaskStatus(**default_task_data)]) failed_isrcs = response.create_error_response( code='', message={ 'territories': [], 'failed_territories': ['AD', 'AF'] }) lock_isrc_result = [('QA123', failed_isrcs)] bulk.save_bulk_lock_report.delay( lock_isrc_results=lock_isrc_result, task_id=1, territories=['AF', 'AD'], reason='test', lock_tasks_amount=1) updated_task = bulk_tasks.get_task(1).message.as_dict() expected_task_result = { 'reason': 'test', 'successful_isrcs': False, 'failed_isrcs': True, 'isrcs': { 'QA123': {'failed_territories': ['AD', 'AF'], 'territories': []}} } assert updated_task['result'] == expected_task_result assert updated_task['status'] == 'FAILED' bulk.generate_reports.assert_called_once_with(1, ['isrc']) def test_update_isrc_for_upc_with_carve_in_territories(): """Should return error because there are no carve in territories.""" isrc_info = {'isrc': 'QA123', 'tuid': '123', 'vendor_id': 7} update_response = bulk.update_isrc_for_upc( upc='UPC123', isrc_info=isrc_info, carve_in_territories=None, user='oa:123', correlation_id='test') expected_response = { 'error_message': 'Unable to find territories carve-in', 'tuid': '123', 'upc': 'UPC123', 'isrc': 'QA123', 'success': False, 'vendor_id': 7 } assert update_response == expected_response @pytest.mark.parametrize( 'upc_statuses, expected_status', [([{'error_count': 0, 'success_count': 1, 'warning_count': 0}], 'DONE'), ([{'error_count': 1, 'success_count': 0, 'warning_count': 0}], 'FAILED'), ([{'error_count': 1, 'success_count': 1, 'warning_count': 0}], 'PARTIAL_FAIL')]) def test_make_upc_report_status( upc_statuses, expected_status, feature_engine): """Should set status correctly.""" report_status = bulk._make_upc_report_status(upc_statuses) assert report_status == expected_status def test_make_upc_report_status_warning(): """Should set warning status.""" upc_statuses = [ {'error_count': 0, 'warning_count': 1, 'success_count': 0}] report_status = bulk._make_upc_report_status(upc_statuses) assert report_status == 'WARNING' def valid_territorries_patch(territories, correlation_id=None): """A quick patch for _is_valid_territories.""" return response.Response(territories) @patch( 'masters_registry.models.ows_territories.convert_territories', new=patches.ows_territories_convert_territories_patch) @patch( 'masters_registry.models.ownership.get_upc_by_tuid', new=patches.ownership_get_upc_by_tuid_patch) @patch( 'masters_registry.models.ows_carveouts.get_dms_carveout_for_upc', new=patches.ows_carveouts_get_dms_carveout_for_upc_not_carved_out_patch) @patch( 'masters_registry.logic.masters_registry._is_valid_territories', new=valid_territorries_patch) @patch( 'masters_registry.models.ownership.get_tracks', new=patches.ownership_get_tracks_patch) @patch('masters_registry.models.yt_ownership.send_message', new=Mock()) @mock_aws def test_bulk_remove_territories( setup_bulk_status_db, setup_masters_active_table, seed_masters_active_table, setup_new_audit_table, mocker): """Functional test for bulk_remove_territories""" from moto.core import patch_client, patch_resource patch_client(outside_dynamodb_client) patch_resource(dynamodb_resource) test_isrc = 'ZW1010600008' ownership_info = { field_const.ISRC: test_isrc, field_const.LOCKED_TERRITORIES: {}, field_const.TERRITORIES: { 'CA': {field_const.TUID: 12345}, 'MX': [{field_const.TUID: 12345}] } } setup_masters_active_table(dynamodb_resource.meta.client) seed_masters_active_table(dynamodb_resource.meta.client, item=ownership_info) setup_new_audit_table(dynamodb_resource.meta.client) account_type = 'vendor' account_id = 123 items = [{ 'isrc': test_isrc, 'tuid': 12345, 'conflict_date': '2018-12-24', 'conflicting_owner': 'UMG', 'territories': ['MX'] }] correlation_id = '123' task_id = bulk_tasks.create_task( correlation_id=correlation_id, user_id=None, task_type=bulk_tasks_const.BULK_REMOVE_TERRITORIES, count=1, context=None, user_name=None, account_id=account_id, account_type=account_type ) conflict_manager_success_response = { 'status': 'ORCHARD_ASSERTED', 'updated_conflicts_amount': 10 } mock_request.post( api_const.OWS_CONFLICT_MANAGER, '/conflicts/status/bulk', status=200, response=conflict_manager_success_response) spy_ows_conflict = mocker.patch.object( ows_conflict_manager, 'bulk_update_conflict_status', wraps=ows_conflict_manager.bulk_update_conflict_status) bulk.bulk_remove_territories.delay( items, account_type, account_id, correlation_id, task_id, 'oa:123') saved_item = ownership.get_ownership(test_isrc).message expected_saved_item = { 'isrc': test_isrc, 'locked_territories': {}, 'territories': { 'CA': {'tuid': Decimal('12345')} }, } saved_item.pop('updated_timestamp') assert saved_item == expected_saved_item expected_audit_records = [ { 'correlation_id': correlation_id, 'user': '{}:{}'.format(account_type, account_id), 'isrc': test_isrc, 'opcode': opcode_const.AUTO_REMOVE, 'territories': {'MX': Decimal('12345')} }, ] audit_records = ownership.get_ownership_audit(test_isrc).message for record in audit_records: record.pop('timestamp') assert audit_records == expected_audit_records yt_ownership.send_message.assert_called_with( test_isrc, {'CA'}, correlation_id) expected_ows_conflict_args = { 'grouped_conflicts_ids': ['12345|2018-12-24|UMG|release'], 'status': field_const.AUTO_RELEASED_VIA_CONFLICT_MANAGER, 'note': '', 'account_type': account_type, 'account_id': account_id, 'correlation_id': correlation_id, 'orchard_user_id': 'oa:123' } spy_ows_conflict.assert_called_with(**expected_ows_conflict_args) created_task = bulk_tasks.get_task(1).message.as_dict() expected_task = { 'create_datetime': ANY, 'correlation_id': correlation_id, 'count': 1, 'user_name': '', 'finish_datetime': ANY, 'user_id': None, 'result': { 'isrcs': [ { 'tuid': 12345, 'isrc': test_isrc, 'success': True, 'conflicting_owner': 'UMG', 'conflict_date': '2018-12-24', 'failed_reason': '', 'territories': ['MX'] } ] }, 'status': 'DONE', 'id': 1, 'type': 'BULK_REMOVE_TERRITORIES', 'account_id': account_id, 'account_type': account_type } assert created_task == expected_task SUCCESS_REPORT = [ { 'tuid': 12345, 'isrc': 'TEST', 'success': True, 'conflicting_owner': 'UMG', 'conflict_date': '2018-12-24', 'failed_reason': '', 'territories': ['MX'] }, { 'tuid': 12345, 'isrc': 'TEST', 'success': True, 'conflicting_owner': 'UMG', 'conflict_date': '2018-12-24', 'failed_reason': '', 'territories': ['MX'] }, ] FAILED_REPORT = [ { 'tuid': 12345, 'isrc': 'TEST', 'success': False, 'conflicting_owner': 'UMG', 'conflict_date': '2018-12-24', 'failed_reason': 'Server error', 'territories': ['MX'] }, { 'tuid': 12345, 'isrc': 'TEST', 'success': False, 'conflicting_owner': 'UMG', 'conflict_date': '2018-12-24', 'failed_reason': 'Server error', 'territories': ['MX'] }, ] PARTIAL_FAILED_REPORT = [ { 'tuid': 12345, 'isrc': 'TEST', 'success': False, 'conflicting_owner': 'UMG', 'conflict_date': '2018-12-24', 'failed_reason': 'Server error', 'territories': ['MX'] }, { 'tuid': 12345, 'isrc': 'TEST', 'success': True, 'conflicting_owner': 'UMG', 'conflict_date': '2018-12-24', 'failed_reason': '', 'territories': ['MX'] }, ] @pytest.mark.parametrize( 'report,expected_status', [ (SUCCESS_REPORT, bulk_tasks_const.DONE_STATUS), (PARTIAL_FAILED_REPORT, bulk_tasks_const.PARTIAL_FAIL_STATUS), (FAILED_REPORT, bulk_tasks_const.FAILED_STATUS) ] ) def test_save_bulk_remove_territories_report( report, expected_status, setup_bulk_status_db, feature_engine): """Test task for saving bulk remove territories report.""" task_id = bulk_tasks.create_task( correlation_id='test', user_id=None, task_type=bulk_tasks_const.BULK_REMOVE_TERRITORIES, count=2, context=None, user_name=None, account_id=333, account_type='vendor' ) bulk.save_bulk_remove_territories_report.delay( report, task_id) updated_task = bulk_tasks.get_task(task_id).message.as_dict() assert updated_task['status'] == expected_status