"""Product logic.""" from typing import Any, Optional import config from lambdacommon.common_config import logger from src.constants.fields import PRODUCT_EXCLUSION_LIST, PRODUCT_SKIP_DATE, PRODUCT_SKIP_LIST from src.logic.types import OpensearchBulkDocument from src.models import ows_product, ows_track from src.utils import document as document_utils from src.utils.parser import parse_spec_file def _get_product_document( release_id: int, event_type: str, doc_type: str ) -> OpensearchBulkDocument | None: json_doc = {} if event_type != "delete": json_doc = ows_product.get_product_document(release_id) if not json_doc: return None if ( json_doc.get("release_id") in PRODUCT_SKIP_LIST or json_doc.get("release_date") == PRODUCT_SKIP_DATE ): return None if release_id in PRODUCT_EXCLUSION_LIST: doc_type = "delete" doc: OpensearchBulkDocument = { "_index": config.PRODUCT_INDEX_NAME, "_op_type": doc_type, "_id": release_id, # no prefix like CS key. "doc": json_doc, # with this it will do upsert and not return NOT_FOUND error. "doc_as_upsert": True, } logger.info(f"{doc_type} product with id {release_id}") return document_utils.prepare_for_upload(parse_spec_file("releases"), doc) def prepare_doc(event: dict[str, Any]) -> Optional[OpensearchBulkDocument]: """Construct Opensearch document for the product corpus. Args: event (dict): event with product data from Maxwell's Returns: dict: Opensearch document ready for upload """ doc_type = "update" event_type = event.get("type") if not event_type or event.get("_metadata"): logger.warning("An event with incompatible format received: %s", event) return None # Prevent the removal of product document if we receive 'delete' event # from 'release_correction' or 'release_approval_queue' tables. # We need to update the product document in this case. if event_type == "delete" and event.get("table") != "releases": event_type = "insert" elif event_type == "delete" and event.get("table") == "releases": doc_type = "delete" event_type = "delete" return _get_product_document(event["data"]["release_id"], event_type, doc_type) def prepare_doc_from_track_additional_isrc( event: dict[str, Any], ) -> Optional[OpensearchBulkDocument]: """Construct Opensearch document for the product corpus based on track related data. Args: event (dict): event with product data from Maxwell's Returns: dict: Opensearch document ready for upload """ if not event.get("data") or not event.get("data", {}).get("track_id"): return None track_id = event.get("data", {}).get("track_id") track_data = ows_track.get_track_by_tuid(track_id) product_id = track_data.get("product_id") if not product_id or not int(product_id): logger.warning("Got bad track data for tuid: %s", track_id) return None return _get_product_document(int(product_id), "update", "update") def prepare_doc_from_product_additional_upc( event: dict[str, Any], ) -> Optional[OpensearchBulkDocument]: """Construct Opensearch document for the product corpus from product UPC data. Args: event (dict): event with product_additional_upc data from Maxwell's Returns: dict: Opensearch document ready for upload """ data = event.get("data") or {} product_id = data.get("product_id") if not product_id or not int(product_id): logger.warning("Got bad product_additional_upc data: %s", event.get("data")) return None return _get_product_document(int(product_id), "update", "update")