"""Lambda hive_text_recognition function module.""" from httpx import HTTPStatusError from kafka_utils.consumer.deserializer.simple_json import JSONDeserializer from kafka_utils.consumer.source.mapping import EventSourceMessage import sentry_sdk from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration import config from src.connectors import hive from src.connectors import ows_assets from src.logic.constants import STATUS_ERROR, STATUS_OK, STATUS_SKIP, STATUS_SUCCESS from src.logic.event_logger import AssetMezzanineCoverartEventLogger from src.logic.event_message import AssetMezzanineCoverartEventMessage if config.SENTRY_DSN: sentry_sdk.init( dsn=config.SENTRY_DSN, environment=config.ENVIRONMENT, integrations=[AwsLambdaIntegration(timeout_warning=True)], ) json_deserializer = JSONDeserializer() lambda_logger = AssetMezzanineCoverartEventLogger(config.app_logger) def run_hive_ocr_task(data: AssetMezzanineCoverartEventMessage): """Run Hive OCR task and store output to ows-assets.""" # get presigned url try: url = ows_assets.get_presigned_url(data.asset_final_id) except HTTPStatusError as e: if ( e.response.status_code == 404 and config.ENVIRONMENT != config.PROD_ENVIRONMENT ): lambda_logger.set_data( status=STATUS_SKIP, result="Failed to get download url: asset not found.", ) return else: raise e block_text = hive.run_task(url) try: ows_assets.post_results(data.asset_final_id, block_text) except HTTPStatusError as e: if e.response.status_code == 409: lambda_logger.set_data( status=STATUS_SKIP, result="Existing Hive result found for asset." ) return else: raise e lambda_logger.set_data( status=STATUS_SUCCESS, result="Saved Hive result to ows-assets." ) def process_event(payload: AssetMezzanineCoverartEventMessage, event_key=None): """Process event from lambda entry point.""" try: lambda_logger.start(event_key=event_key, payload=payload) run_hive_ocr_task(payload) except Exception as e: lambda_logger.set_data(status=STATUS_ERROR, result=str(e)) raise e finally: lambda_logger.end() def handler(event, context): """Lambda Entry point.""" if event.get("eventSource") == "custom": for asset in event.get("assets", []): data = AssetMezzanineCoverartEventMessage.from_dict(asset) process_event(data) else: msk_message = EventSourceMessage(event) for event_key, message in msk_message: data = AssetMezzanineCoverartEventMessage( message=message.value, topic=message.topic, value_deserializer=json_deserializer, ) process_event(data, event_key=event_key) return {"status": STATUS_OK}