import json import pytest import config import constants env = config.ENV LAMBDA_NAME = constants.LAMBDA_VECTOR_THROTTLER STORE_ID = 465 DELIVERY_QUEUE_NAME_PATTERN = \ '{env}-delivery{encoder_id}_e{priority:07d}_d{dms_priority:07d}' @pytest.mark.parametrize('alchemy_db', ['dd'], indirect=True) def test_lambda_vector_throttler(alchemy_db, lambda_client, sqs_client, cloudwatch_log_handler): """Test lambda_vector_throttler lambda finds and queues orders that are ready to deliver \n https://www.notion.so/Throttler- c4c1aaf4b2f34d87b4c1ca9c0e783db8?pvs=4#394f5418f2754d04af0b40047702bfab""" # update detail queue to have entity for process encoding_queue_detail = \ alchemy_db.encoding_queue.get_encoding_queue_detail_for_deliver( STORE_ID, ["new", "ready_to_encode", "encoding", "encoded", "delivered"], 1, )[0] alchemy_db.encoding_queue.update_queue_detail_status( encoding_queue_detail.encoding_queue_detail_id, 'ready_to_deliver') # check there is at least 1 encoding queue ready for deliver is_any_ready_for_deliver_encoding_queue_detail(alchemy_db, True) # Cleanup SQS before invoking lambda sqs_name = get_sqs_delivery_name(encoding_queue_detail) sqs_client.purge(sqs_name) # Invoke lambda lambda_event = { 'dms_id': STORE_ID, 'limit': 1000, 'allowed_encoders': [18]} lambda_client.invoke(LAMBDA_NAME, lambda_event) cloudwatch_log_handler.assert_lambda_logs_message( LAMBDA_NAME, f'Finished .*lambda-vector-throttler.* for dms_id {STORE_ID}' ) # check encoding queue detail status is updated properly alchemy_db.encoding_queue.assert_encoding_queue_detail_status( encoding_queue_detail.encoding_queue_detail_id, 'queued_for_delivery') # check there are no encoding queue ready for deliver is_any_ready_for_deliver_encoding_queue_detail(alchemy_db, False) # # Check SQS assert_sqs_message_data(sqs_client, sqs_name, encoding_queue_detail) def is_any_ready_for_deliver_encoding_queue_detail(db, exists): encoding_queues_for_deliver = \ db.encoding_queue.get_encoding_queue_detail_for_deliver( STORE_ID, ['ready_to_deliver']) if exists: assert len(encoding_queues_for_deliver) > 0 else: assert len(encoding_queues_for_deliver) == 0 def get_sqs_delivery_name(encoding_queue): encoder_id = encoding_queue.encoder_id priority = encoding_queue.encoding_order_priority sqs_name = DELIVERY_QUEUE_NAME_PATTERN.format( env=env, encoder_id=encoder_id, priority=int(priority), dms_priority=1 ) return sqs_name def assert_sqs_message_data(sqs_client, sqs_name, encoding_queue): message = sqs_client.read_message(sqs_name) message_body = json.loads(message['Body']) assert message_body['dms_id'] == STORE_ID assert message_body['eqd_id'] == encoding_queue.encoding_queue_detail_id