"""Logic for getting indexed products and sending to SQS.""" import json from src import config from src.connectors import get_os_connector from src.connectors import get_sqs_connector # this method is also used in reindex-product-metadata def get_all_indexes(): """Get all items in the review queue index.""" total_count = get_os_connector().os.count(index=config.OPENSEARCH_INDEX_ALIAS).get('count') config.app_logger.debug( 'Fetching {} total products from OpenSearch index.'.format(total_count)) search_response = get_os_connector().os.search( index=config.OPENSEARCH_INDEX_ALIAS, _source=True, _source_includes=['review_queue_id', 'product_id', 'queue_name'], size=total_count ).get('hits').get('hits') included_queue_names = ['initial', 'under_investigation', 'escalation'] return [{ 'reviewQueueId': item.get('_source').get('review_queue_id'), 'productId': item.get('_source').get('product_id'), 'queueName': item.get('_source').get('queue_name'), } for item in search_response if item.get('_source').get('queue_name') in included_queue_names] def sqs_send_message(indexes): """Send product ids to SQS queue.""" success_count = 0 for ids in indexes: product_id = ids['productId'] msg_body = { 'operation': config.UPDATE, 'product_id': product_id, 'review_queue_id': ids['reviewQueueId'] } response = get_sqs_connector().send_message( QueueUrl=config.SQS_QUEUE_URL, MessageBody=json.dumps(msg_body) ) status = response.get('ResponseMetadata').get('HTTPStatusCode') if status == 200: success_count += 1 else: config.app_logger.warning(f'{status} error adding product {product_id} to SQS queue') config.app_logger.info(f'{success_count} Product IDs were successfully added to the SQS queue') return {'status': 'complete'}