"""Lambda function for submitting documents to CloudSearch.""" import base64 import json import logging from aws_kinesis_agg import deaggregator import boto3 import sentry_sdk import config from logic import artist_info from logic import product from logic import project from logic import subaccount from logic import track from logic import users from logic import vendor import sentry logger = logging.getLogger(config.APPLICATION_NAME) sentry.init_client() RELEASE_TABLES = ('releases', 'release_correction', 'release_approval_queue') USER_TABLES = ('vend_contact', 'contact', 'subaccount', 'vend_contact_roles') VENDOR_TABLES = ('vendor', 'vend_contact') def handler(event, context): """Lambda entry point.""" artists_documents = [] track_documents = [] project_documents = [] releases_documents = [] release_tracks_documents = [] subaccount_documents = [] vendor_documents = [] user_documents = [] doc = None logger.info(event) for raw_record in deaggregator.iter_deaggregate_records(event['Records']): record_str = base64.b64decode(raw_record['kinesis']['data']) record_json = json.loads(record_str) # branch on table if record_json.get('table') == 'track': doc = track.prepare_doc(record_json) if doc: track_documents.append(doc) if record_json.get('table') == 'project': doc = project.prepare_doc(record_json) if doc: project_documents.append(doc) if record_json.get('table') in RELEASE_TABLES: doc = product.prepare_doc(record_json) if doc: releases_documents.append(doc) if record_json['table'] == 'releases' and record_json['type'] == 'update': tracks_docs = track.prepare_release_tracks_docs(record_json) if tracks_docs: release_tracks_documents.extend(tracks_docs) if record_json.get('table') == 'subaccount': doc = subaccount.prepare_doc(record_json) if doc: subaccount_documents.append(doc) if record_json.get('table') in VENDOR_TABLES: doc = vendor.prepare_doc(record_json) if doc: vendor_documents.append(doc) if record_json.get('table') in USER_TABLES: docs = users.prepare_docs(record_json) if docs: user_documents = user_documents + docs if record_json.get('table') == 'artist_info': doc = artist_info.prepare_doc(record_json) if doc: artists_documents.append(doc) if track_documents: logger.info(f'Preparing to upload {len(track_documents)} track documents') data = json.dumps(track_documents) upload_documents('tracks', data) if project_documents: logger.info(f'Preparing to upload {len(project_documents)} project documents') data = json.dumps(project_documents) upload_documents('projects', data) if releases_documents: logger.info(f'Preparing to upload {len(releases_documents)} releases documents') release_data = json.dumps(releases_documents) upload_documents('releases', release_data) if release_tracks_documents: logger.info(f'Preparing to upload {len(release_tracks_documents)} release tracks documents') tracks_data = json.dumps(release_tracks_documents) upload_documents('tracks', tracks_data) if subaccount_documents: logger.info(f'Preparing to upload {len(subaccount_documents)} subaccount documents') data = json.dumps(subaccount_documents) upload_documents('labels', data) if vendor_documents: logger.info(f'Preparing to upload {len(vendor_documents)} vendor documents') data = json.dumps(vendor_documents) upload_documents('labels', data) if user_documents: logger.info(f'Preparing to upload {len(user_documents)} user documents') data = json.dumps(user_documents) upload_documents('users', data) if artists_documents: logger.info(f'Preparing to upload {len(artists_documents)} artists documents') data = json.dumps(artists_documents) upload_documents('artists', data) return True def upload_documents(corpus, documents): """Upload documents to CloudSearch corpus. Args: corpus: (str) tracks|releases|projects documents: (array) array of documents for a corpus """ aws_client = boto3.client( 'cloudsearchdomain', endpoint_url=config.cloudsearch_endpoints[ corpus]) logger.info(corpus) logger.info(documents) try: response = aws_client.upload_documents( documents=documents, contentType='application/json') logger.info(response) except Exception as upload_exception: # TODO: send a document to DLQ and do not raise an exception # for CloudSearch upload errors logger.exception(upload_exception) sentry_sdk.capture_exception() raise