"""Module to handle messaging logic.""" import uuid from typing import Any from flask import g from owsresponse import response from transcoding import config from transcoding.connectors import ows_assets from transcoding.connectors.ows_assets import ( AssetOwnerNotDeterminableException, ) from transcoding.models import ( sqs_message, transcoding_job as transcoding_job_model, transcoding_order as transcoding_order_model, ) def get_message_group_id(input_key: str) -> str: """Get message group id for SQS queue. Args: input_key (str): Input key of transcoding order. Returns: str: Message group id. """ try: asset_upload = ows_assets.get_asset_owner(input_key) return str(asset_upload["vendor_id"]) + "-" + str(asset_upload["subaccount_id"]) except AssetOwnerNotDeterminableException: g.log.warning( "Unable to get asset owner for input_key=%s", input_key, exc_info=True ) return str(uuid.uuid4()) except Exception: g.log.exception("Unable to get asset owner for input_key=%s", input_key) return str(uuid.uuid4()) def send_transcoding_job_message( transcoding_order: dict[str, Any], transcoding_job: dict[str, Any] ) -> response.Response: """Send transcoding job message. Args: transcoding_order (dict): Transcoding order instance. transcoding_job (dict): Transcoding job instance. Returns: response.Response: Message with success or error message. """ message_group_id = get_message_group_id(transcoding_order["input_key"]) if not transcoding_job["codec"]: message = { "transcoding_job_id": transcoding_job["transcoding_job_id"], "input_bucket": transcoding_order["input_bucket"], "input_key": transcoding_order["input_key"], "output_bucket": transcoding_job["output_bucket"], "output_key": transcoding_job["output_key"], "pass_thru": True, } else: message = { "transcoding_job_id": transcoding_job["transcoding_job_id"], "input_bucket": transcoding_order["input_bucket"], "input_key": transcoding_order["input_key"], "output_bucket": transcoding_job["output_bucket"], "output_key": transcoding_job["output_key"], "container": transcoding_job["container"], "channels": transcoding_job["channels"], "codec": transcoding_job["codec"], "sample_rate": transcoding_job["sample_rate"], "bit_rate": transcoding_job["bit_rate"], "bit_depth": transcoding_job["bit_depth"], } return sqs_message.send_sqs_message( queue_url=config.SQS_TRANSCODING_URL, message=message, message_group_id=message_group_id, ) def resend_transcoding_job_sqs_message(transcoding_job_id: int) -> response.Response: """Resend sqs message for given transcoding_job_id. Args: transcoding_job_id (int): Transcoding job id. Returns: response.Response: Success or error message. """ transcoding_job_response = transcoding_job_model.get_transcoding_job_by_id( transcoding_job_id ) if not transcoding_job_response: return transcoding_job_response transcoding_job = transcoding_job_response.message transcoding_order_id = transcoding_job["transcoding_order_id"] transcoding_order_response = transcoding_order_model.get_transcoding_order( transcoding_order_id ) if not transcoding_order_response: return transcoding_order_response transcoding_order = transcoding_order_response.message send_job_response = send_transcoding_job_message(transcoding_order, transcoding_job) return send_job_response