"""Elasticsearch connector.""" from functools import cache from content_utils.connectors.opensearch import LambdaOpensearchConnector from src import config @cache def get_os_connector_write(): """Get OS Connector.""" os_opts = {} if config.OPENSEARCH_PORT: os_opts.update(port=config.OPENSEARCH_PORT) if config.OPENSEARCH_INDEX_ALIAS_WRITE: os_opts.update(alias_name=config.OPENSEARCH_INDEX_ALIAS_WRITE) return LambdaOpensearchConnector( config.OPENSEARCH_ENDPOINT, config.app_logger, **os_opts ) @cache def get_os_connector_read(): """Get OS Connector.""" os_opts = {} if config.OPENSEARCH_PORT: os_opts.update(port=config.OPENSEARCH_PORT) if config.OPENSEARCH_INDEX_ALIAS_READ: os_opts.update(alias_name=config.OPENSEARCH_INDEX_ALIAS_READ) return LambdaOpensearchConnector( config.OPENSEARCH_ENDPOINT, config.app_logger, **os_opts ) def drop_index(): """Drop the review queue index.""" drop_response = get_os_connector_write().clear_index_via_alias() config.app_logger.debug(drop_response) def get_all_items(): """Get all items in the review queue index.""" total_count = get_os_connector_read().os.count( index=config.OPENSEARCH_INDEX_ALIAS_READ).get('count') config.app_logger.debug(f'Fetching {total_count} total products from index.') search_response = get_os_connector_read().os.search( index=config.OPENSEARCH_INDEX_ALIAS_READ, _source=True, _source_includes=['review_queue_id', 'product_id'], size=total_count ).get('hits').get('hits') return [{ 'reviewQueueId': item.get('_source').get('review_queue_id'), 'productId': item.get('_source').get('product_id') } for item in search_response]