"""Test Queues work.""" from queue import Queue from time import sleep from unittest import mock from flexmock import flexmock from moto import mock_sqs import pytest from transcoding.connectors import sqs from transcoding.connectors.logger import app_logger from transcoding.constants.exceptions import FailedToStartError from transcoding.logic import processing from transcoding.logic import queues from tests.test_utils import sqs_utils PREFETCH_NUMBER = 2 VISIBILITY_TIMEOUT = 60 WAIT_TIME_SECONDS = 2 NUM_MESSAGES = 4 RETURN_TO_QUEUE_TIMEOUT = 1 MAX_RETRIES = 10 STOP_TIMEOUT = 5 PREFETCH_QUEUE_GET_TIMEOUT = 5 SQS_QUEUE_NAME = 'test-q' KEY_FILE = 'key' READERS_COUNT = 2 EXECUTORS_COUNT = 2 @pytest.fixture(scope='session', autouse=True) def start_sqs_mock(request): """Use this function for SQS mock.""" mock = mock_sqs() mock.start() def start_sqs_mock_teardown(): mock.stop() request.addfinalizer(start_sqs_mock_teardown) @pytest.fixture def prefetch_queue(): """Prefetch queue with this fucntion.""" pref_queue = Queue(PREFETCH_NUMBER) return pref_queue @pytest.fixture() def sqs_queue(request): """Connect to SQS queue.""" sqs_resource = sqs.get_sqs_resources() queue_mocked = sqs_resource.create_queue(QueueName='string') mock_queue = sqs.get_queue(queue_mocked.url, sqs_resource=sqs_resource) def sqs_queue_teardown(): mock_queue.delete() request.addfinalizer(sqs_queue_teardown) return mock_queue @pytest.fixture def queue_reader_worker(request, prefetch_queue, sqs_queue): """Create queue reader worker.""" worker = queues.QueueReaderWorker( 'test-worker', prefetch_queue, sqs_queue, VISIBILITY_TIMEOUT, WAIT_TIME_SECONDS, NUM_MESSAGES, RETURN_TO_QUEUE_TIMEOUT, MAX_RETRIES, app_logger ) worker.debug = True def queue_reader_worker_teardown(): worker.stop(STOP_TIMEOUT) request.addfinalizer(queue_reader_worker_teardown) return worker @pytest.fixture def queue_reader_worker_mocked_sqs(request, prefetch_queue): """Create queue reader worker.""" worker = queues.QueueReaderWorker( 'test-worker', prefetch_queue, mock.MagicMock(), VISIBILITY_TIMEOUT, WAIT_TIME_SECONDS, NUM_MESSAGES, RETURN_TO_QUEUE_TIMEOUT, MAX_RETRIES, app_logger ) worker.debug = True def queue_reader_worker_teardown(): worker.stop(STOP_TIMEOUT) request.addfinalizer(queue_reader_worker_teardown) return worker @pytest.fixture def app_execution_worker(request, prefetch_queue, sqs_queue): """Create queue executive worker.""" worker_id = 1 worker = queues.AppProcessingWorker( 'test-worker', prefetch_queue, sqs_queue, RETURN_TO_QUEUE_TIMEOUT, app_logger, worker_id ) worker.debug = True def api_execution_worker_teardown(): worker.stop(STOP_TIMEOUT) request.addfinalizer(api_execution_worker_teardown) return worker @pytest.fixture def worker_manager(): """Create queue worker manager.""" manager = queues.QueueWorkerManager( READERS_COUNT, EXECUTORS_COUNT, SQS_QUEUE_NAME, PREFETCH_NUMBER, VISIBILITY_TIMEOUT, WAIT_TIME_SECONDS, NUM_MESSAGES, app_logger ) manager.debug = True return manager @pytest.fixture def test_connection(): """Make mock connection.""" sqs_resource = sqs.get_sqs_resources() flexmock(sqs).should_receive('get_sqs_resources').and_return(sqs_resource) return sqs_resource def test_read_message(queue_reader_worker, prefetch_queue, sqs_queue): """Verify we can read all messages, sent to queue.""" # test functional calls sqs_utils.put_sqs_messages(sqs_queue, PREFETCH_NUMBER) queue_reader_worker.start() counter = 0 while counter < PREFETCH_NUMBER: prefetch_queue.get(timeout=PREFETCH_QUEUE_GET_TIMEOUT) counter += 1 # checking assert counter == PREFETCH_NUMBER def test_read_message_prefetch_queue_full( queue_reader_worker, prefetch_queue, sqs_queue): """Test all messages readability after full prefetch. Verify we can read all messages, sent to queue, assuming even if the prefetch queue becomes full """ # test functional calls sqs_utils.put_sqs_messages(sqs_queue, PREFETCH_NUMBER * 2) queue_reader_worker.start() counter = 0 while counter < PREFETCH_NUMBER * 2: prefetch_queue.get(timeout=PREFETCH_QUEUE_GET_TIMEOUT) counter += 1 # checking assert counter == PREFETCH_NUMBER * 2 @mock.patch('moto.sqs.models.Message.change_visibility') def test_read_message_prefetch_queue_full_error( change_visibility_mock, queue_reader_worker, sqs_queue): """Verify message deletion from SQS after prefetch. Verify we return the message back to the SQS queue, if prefetch queue becomes full during the prefetch process """ # test functional calls sqs_utils.put_sqs_messages(sqs_queue, 2) queue_reader_worker.prefetch_queue = Queue(1) queue_reader_worker.num_messages = 2 queue_reader_worker.start() sleep(3) # checking assert change_visibility_mock.called change_visibility_mock.assert_called_with(VISIBILITY_TIMEOUT) def test_read_message_sqs_delete_success( queue_reader_worker_mocked_sqs): """Verify message deletion from SQS after prefetch. Verify we delete message from SQS queue, if it was passed to processing layer. """ # mocking message_1 = { 'MessageId': 'message_id_1', 'ReceiptHandle': 'receipt_handle_1'} message_2 = { 'MessageId': 'message_id_2', 'ReceiptHandle': 'receipt_handle_2'} message_1_mock = mock.MagicMock() message_1_mock.message_id = message_1['MessageId'] message_1_mock.receipt_handle = message_1['ReceiptHandle'] message_2_mock = mock.MagicMock() message_2_mock.message_id = message_2['MessageId'] message_2_mock.receipt_handle = message_2['ReceiptHandle'] mock_messages = [ message_1_mock, message_2_mock] sqs_queue_mock = mock.MagicMock() receive_messages_mock = mock.MagicMock( return_value=mock_messages) delete_messages_mock = mock.MagicMock(return_value={}) sqs_queue_mock.delete_messages = delete_messages_mock sqs_queue_mock.receive_messages = receive_messages_mock queue_reader_worker_mocked_sqs.sqs_queue = sqs_queue_mock queue_reader_worker_mocked_sqs.start() delete_entries = [ { 'Id': 'message_id_1', 'ReceiptHandle': 'receipt_handle_1' }, { 'Id': 'message_id_2', 'ReceiptHandle': 'receipt_handle_2' }] sleep(3) queue_reader_worker_mocked_sqs.stop(timeout=0) receive_messages_mock.assert_called_once() delete_messages_mock.assert_called_once_with(Entries=delete_entries) def test_read_message_sqs_delete_failed(queue_reader_worker_mocked_sqs): """Verify sentry call if message deletion from SQS failed. Verify we notify sentry if delete message from SQS queue failed. """ # mocking message_1 = { 'MessageId': 'message_id_1', 'ReceiptHandle': 'receipt_handle_1'} message_2 = { 'MessageId': 'message_id_2', 'ReceiptHandle': 'receipt_handle_2'} message_1_mock = mock.MagicMock() message_1_mock.message_id = message_1['MessageId'] message_1_mock.receipt_handle = message_1['ReceiptHandle'] message_2_mock = mock.MagicMock() message_2_mock.message_id = message_2['MessageId'] message_2_mock.receipt_handle = message_2['ReceiptHandle'] mock_messages = [ message_1_mock, message_2_mock] sqs_queue_mock = mock.MagicMock() receive_messages_mock = mock.MagicMock( return_value=mock_messages) delete_messages_mock = mock.MagicMock(return_value={ 'Successful': [ { 'Id': message_1['MessageId'] }], 'Failed': [ { 'Id': message_2['MessageId'], 'SenderFault': False, 'Code': 'error_code', 'Message': 'error_message' }]}) sqs_queue_mock.delete_messages = delete_messages_mock sqs_queue_mock.receive_messages = receive_messages_mock queue_reader_worker_mocked_sqs.sqs_queue = sqs_queue_mock queue_reader_worker_mocked_sqs.start() delete_entries = [ { 'Id': 'message_id_1', 'ReceiptHandle': 'receipt_handle_1' }, { 'Id': 'message_id_2', 'ReceiptHandle': 'receipt_handle_2' }] sleep(3) receive_messages_mock.assert_called_once() delete_messages_mock.assert_called_once_with(Entries=delete_entries) queue_reader_worker_mocked_sqs.stop(timeout=0) def test_read_message_sqs_delete_queue_full( queue_reader_worker_mocked_sqs): """Verify message deletion from SQS after prefetch. Verify we delete message from SQS queue, if it was passed to processing layer. """ # mocking message_1 = { 'MessageId': 'message_id_1', 'ReceiptHandle': 'receipt_handle_1'} message_2 = { 'MessageId': 'message_id_2', 'ReceiptHandle': 'receipt_handle_2'} message_1_mock = mock.MagicMock() message_1_mock.message_id = message_1['MessageId'] message_1_mock.receipt_handle = message_1['ReceiptHandle'] message_2_mock = mock.MagicMock() message_2_mock.message_id = message_2['MessageId'] message_2_mock.receipt_handle = message_2['ReceiptHandle'] mock_messages = [ message_1_mock, message_2_mock] sqs_queue_mock = mock.MagicMock() receive_messages_mock = mock.MagicMock( return_value=mock_messages) delete_messages_mock = mock.MagicMock(return_value={}) sqs_queue_mock.delete_messages = delete_messages_mock sqs_queue_mock.receive_messages = receive_messages_mock queue_reader_worker_mocked_sqs.sqs_queue = sqs_queue_mock queue_reader_worker_mocked_sqs.prefetch_queue = Queue(1) queue_reader_worker_mocked_sqs.num_messages = 2 queue_reader_worker_mocked_sqs.start() delete_entries = [ { 'Id': 'message_id_1', 'ReceiptHandle': 'receipt_handle_1' }] sleep(3) receive_messages_mock.assert_called_once() delete_messages_mock.assert_called_once_with(Entries=delete_entries) queue_reader_worker_mocked_sqs.stop(timeout=0) def test_start(worker_manager, test_connection): """Verify that the manager can start a defined number of workers.""" # mocking mocked_queue = test_connection.create_queue(QueueName=SQS_QUEUE_NAME) # test functional call worker_manager.start() # checking assert worker_manager.readers_started == READERS_COUNT assert worker_manager.executors_started == EXECUTORS_COUNT worker_manager.stop(STOP_TIMEOUT) mocked_queue.delete() def test_start_failed(worker_manager): """Verify FailedToStartError is raised properly. Verify that FailedToStartError is raised if no workers were able to start successfully. """ # mocking (flexmock(sqs).should_receive('get_queue') .with_args(SQS_QUEUE_NAME) .and_return(None)) # checking with pytest.raises(FailedToStartError): # test functional call worker_manager.start() def test_process_message_success(prefetch_queue, app_execution_worker): """Verify that we delete the message from SQS after processing it.""" message = mock.MagicMock() worker_id = 1 (flexmock(processing) .should_receive('process_transcoding_job') .with_args(message, worker_id) .and_return((True, '')) .once()) prefetch_queue.put(message) app_execution_worker.start() prefetch_queue.join() assert message.delete.called