"""Test Queues work.""" from queue import Queue from time import sleep from unittest import mock import flexmock from moto import mock_sqs from oto import response import pytest from app_name.connectors import sqs from app_name.connectors.loggly import app_logger from app_name.constants.exceptions import FailedToStartError from app_name.logic import queue from app_name.logic import processing_sample 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 = queue.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 app_execution_worker(request, prefetch_queue, sqs_queue): """Create queue executive worker.""" worker = queue.AppProcessingWorker( 'test-worker', prefetch_queue, sqs_queue, RETURN_TO_QUEUE_TIMEOUT, app_logger ) 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 = queue.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_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() (flexmock(processing_sample) .should_receive('sample_processing') .with_args(message) .and_return(response.Response())) prefetch_queue.put(message) app_execution_worker.start() prefetch_queue.join() assert message.delete.called