"""Lambda ar-search-ingest function module.""" import base64 import json from typing import Any, List import sentry_sdk from aws_kinesis_agg import deaggregator from aws_lambda_powertools.utilities.data_classes import KinesisStreamEvent from lambdacommon import util from lambdacommon.common_config import logger from src.constants.table_mapping import ( PRODUCT_ADDITIONAL_UPC_TABLES, PROJECT_TABLES, RELEASE_TABLES, TRACK_ADDITIONAL_ISRC_TABLES, TRACK_TABLES, ) from src.logic import product, project, track from src.logic.common import is_duplicate_event from src.logic.types import OpensearchBulkDocument from src.models.opensearch import upload_bulk_documents util.init_sentry_for_lambda() def handler(event: KinesisStreamEvent, context: Any) -> dict[str, str]: """Lambda entry point.""" try: logger.debug(event) if not event.get("Records"): return {"status": "BAD_INPUT"} documents: List[OpensearchBulkDocument] = [] for raw_record in deaggregator.iter_deaggregate_records(event["Records"]): record_str = base64.b64decode(raw_record["kinesis"]["data"]) record_json = json.loads(record_str) logger.debug(record_json) # branch on table names. doc = None if is_duplicate_event(record_json, documents): logger.info(f"Duplicate event: {json.dumps(record_json)}") continue if record_json.get("table") in RELEASE_TABLES: doc = product.prepare_doc(record_json) if record_json["table"] == "releases" and record_json["type"] == "update": tracks_docs = track.prepare_release_tracks_docs(record_json) if tracks_docs: documents.extend(tracks_docs) elif record_json.get("table") in PROJECT_TABLES: doc = project.prepare_doc(record_json) elif record_json.get("table") in TRACK_TABLES: doc = track.prepare_doc(record_json) elif record_json.get("table") in TRACK_ADDITIONAL_ISRC_TABLES: doc = product.prepare_doc_from_track_additional_isrc(record_json) elif record_json.get("table") in PRODUCT_ADDITIONAL_UPC_TABLES: doc = product.prepare_doc_from_product_additional_upc(record_json) else: logger.info(f"SKIP: No mapping for table: {record_json.get('table')}") # No need keep data separate per index as index_name is part of data if doc: documents.append(doc) # We are not doing batching here instead rely on input batching to be # small enough not to avoid overwhelming memory and cluster. upload_bulk_documents(documents) return {"status": "OK"} except Exception as e: logger.exception(str(e)) sentry_sdk.capture_exception() sentry_sdk.flush() raise e