from queue import Queue from flexmock import flexmock from moto import mock_aws import pytest from tests.utils import sqs as sqs_utils from ytownership.connectors import sqs from ytownership.connectors import youtube from ytownership.connectors.loggly import app_logger from ytownership.connectors.sqs import JSONMessageExt from ytownership.logic import message as message_logic from ytownership.logic import queue from ytownership.utils.misc import DotDict PREFETCH_NUMBER = 2 VISIBILITY_TIMEOUT = 60 WAIT_TIME_SECONDS = 2 NUM_MESSAGES = 4 RETURN_TO_QUEUE_TIMEOUT = 1 MAX_RETRIES = 3 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): mock = mock_aws() mock.start() def start_sqs_mock_teardown(): mock.stop() request.addfinalizer(start_sqs_mock_teardown) @pytest.fixture def prefetch_queue(): pref_queue = Queue(PREFETCH_NUMBER) return pref_queue @pytest.fixture() def sqs_queue(request): conn = sqs.get_connection() conn.create_queue(QueueName=SQS_QUEUE_NAME) mock_queue = sqs.get_queue(SQS_QUEUE_NAME, sqs_connection=conn) 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): 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 api_execution_worker(request, prefetch_queue, sqs_queue): worker = queue.ApiExecutorWorker( 'test-worker', prefetch_queue, sqs_queue, object(), 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(): mock_service_mapping = DotDict({ 'youtube_partner': object() }) (flexmock(youtube).should_receive('ServiceMapping') .and_return(mock_service_mapping)) manager = queue.QueueWorkerManager( READERS_COUNT, EXECUTORS_COUNT, SQS_QUEUE_NAME, PREFETCH_NUMBER, VISIBILITY_TIMEOUT, WAIT_TIME_SECONDS, NUM_MESSAGES, KEY_FILE, app_logger ) manager.debug = True return manager @pytest.fixture @mock_aws def test_connection(): test_conn = sqs.get_connection() flexmock(sqs).should_receive('get_connection').and_return(test_conn) return test_conn def test_read_message(queue_reader_worker, prefetch_queue, sqs_queue): """Verify we can read all messages, sent to queue """ 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 assert counter == PREFETCH_NUMBER def test_read_message_prefetch_queue_full( queue_reader_worker, prefetch_queue, sqs_queue): """Verify we can read all messages, sent to queue, assuming even if the prefetch queue becomes full """ 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 assert counter == PREFETCH_NUMBER * 2 def test_process_message_success(prefetch_queue, api_execution_worker): """Verify that we delete the message from SQS after processing it successfully """ message = sqs_utils.create_sqs_message() prefetch_queue.put(JSONMessageExt(message)) flexmock(message_logic).should_receive('process_message').and_return( message_logic.MessageProcessingResult.SUCCESS ) flexmock(JSONMessageExt).should_receive('delete').at_least().once() api_execution_worker.start() prefetch_queue.join() def test_process_message_failure(prefetch_queue, api_execution_worker): """Verify that we delete the message from SQS after processing failure """ message = sqs_utils.create_sqs_message() prefetch_queue.put(JSONMessageExt(message)) flexmock(message_logic).should_receive('process_message').and_return( message_logic.MessageProcessingResult.FAILURE ) flexmock(JSONMessageExt).should_receive('delete').at_least().once() api_execution_worker.start() prefetch_queue.join() def test_process_message_error(prefetch_queue, api_execution_worker): """Verify that the message is being returned to SQS queue in case of recoverable error """ (flexmock(JSONMessageExt). should_receive('change_visibility').at_least().once()) flexmock(message_logic).should_receive('process_message').and_return( message_logic.MessageProcessingResult.ERROR ) message = sqs_utils.create_sqs_message() prefetch_queue.put(JSONMessageExt(message)) api_execution_worker.start() prefetch_queue.join() def test_start(worker_manager, test_connection): """Verify that the manager can start a defined number of workers """ test_connection.create_queue(QueueName=SQS_QUEUE_NAME) worker_manager.start() assert worker_manager.readers_started == READERS_COUNT assert worker_manager.executors_started == EXECUTORS_COUNT worker_manager.stop(STOP_TIMEOUT) sqs.get_queue(SQS_QUEUE_NAME).delete()