"""Module for processing of incoming job messages.""" from sentry_sdk import capture_exception from transcoding.connectors import logger from transcoding.constants import exceptions as transcoding_exceptions from transcoding.logic import job_status from transcoding.logic import metadata_capture from transcoding.logic import storage from transcoding.logic import transcoding_logic from transcoding.utils import message_extractor from transcoding.validation import validate_job def process_transcoding_job(message, worker_id): """Process transcoding job. Args: message (sqs.Message): message pulled from SQS-queue. worker_id (int): Worker id. Returns: (bool, str): Tuple of result and description. """ current_app_logger = logger.get_current_logger() transcoding_job_id = None try: message_body = message_extractor.get_message_body(message) if not isinstance(message_body, dict): raise transcoding_exceptions.TranscodingFatalError( 'Message body is not a valid JSON object.' ) transcoding_job_id = message_body.get('transcoding_job_id') pass_through = message_body.get('pass_thru', False) validate_job.validate_transcoding_job_schema(message_body, pass_through) output_bucket = message_body['output_bucket'] output_key = message_body['output_key'] input_bucket = message_body['input_bucket'] input_key = message_body['input_key'] current_app_logger.info( f'processing_transcoding_job : {transcoding_job_id},{output_bucket}/{output_key},{input_bucket}/{input_key}' ) if pass_through: storage.copy_object(input_bucket, input_key, output_bucket, output_key) presigned_output_url = storage.generate_url_for_uploaded_asset(output_bucket, output_key) asset_metadata = metadata_capture.get_metadata(presigned_output_url) else: asset_output_path = transcoding_logic.transcode_asset( message, worker_id) storage.upload_job_results(asset_output_path, message_body) asset_metadata = metadata_capture.get_metadata(asset_output_path) job_status.post_completed_status( transcoding_job_id, asset_metadata, ) return True, '' except transcoding_exceptions.TranscodingFatalError as e: capture_exception(e) job_status.post_error_status( transcoding_job_id, str(e)) return False, str(e) except transcoding_exceptions.TranscodingRetryableError as e: capture_exception(e) job_status.post_error_status( transcoding_job_id, str(e), retryable=True) return False, str(e) except Exception as e: error = 'Unexpected processing error: {error}'.format(error=str(e)) capture_exception(e) return False, error