"""Transfer from s3 to storage.""" import json import os from pathlib import Path from video import config from video.connectors import s3 from video.constants import assets from video.constants import job_io_fields from video.constants import job_statuses from video.logic.activity_task_logger import activity_task_logger from video.models import ows_video ASSET_TYPE_PRORES = 'PRORES' ASSET_TYPE_H264 = 'H264' ASSET_TYPE_TIFF_S2 = 'TIFF_S2' ASSET_TYPE_TIFF_S5 = 'TIFF_S5' def get_ripper_h264_filename(inputs): """Get the correct template to use by video resolution. Args: inputs (dict): Inputs. Returns: template (string): The H264 filename template for ripper. """ mezzanine_metadata = json.loads( inputs[job_io_fields.PRORES_MEZZANINE_METADATA]) track_video_resolution = ows_video.get_track_video_resolution( mezzanine_metadata) template = assets.RESOLUTION_TO_RIPPER_TEMPLATE_MAP[track_video_resolution] return template.format(upc=inputs[job_io_fields.UPC]) @activity_task_logger([ job_io_fields.PRORES_MEZZANINE_OUTPUT_S3_KEY, job_io_fields.UPC ], is_last_handler=False, is_first_handler=True) def transfer_prores_from_s3_to_ripper(inputs): """Transfer Prores file from s3 to ripper windows share. Args: inputs (dict): Inputs. Returns: dict: empty. """ prores_mezzanine_path = ( inputs[job_io_fields.PRORES_MEZZANINE_OUTPUT_S3_KEY]) video_s3_bucket = assets.VIDEO_S3_BUCKET upc = inputs[job_io_fields.UPC] filename = assets.PRORES_FILENAME_RIPPER_TEMPLATE.format(upc=upc) on_prem_prores_mezzanine_path = os.path.join( config.RIPPER_OUTPUT_VIDEO_WINDOWS_MOUNT_PATH, config.OUTPUT_VIDEO_PRORES_MEZZANINE_DIR, str(upc)) try: os.mkdir(on_prem_prores_mezzanine_path) except FileExistsError: pass windows_storage_path = os.path.join( on_prem_prores_mezzanine_path, filename) s3_client = s3.get_s3_client() s3_client.download_file( video_s3_bucket, prores_mezzanine_path, windows_storage_path) done_file_name = '{}.done'.format(upc) done_file_path = os.path.join( on_prem_prores_mezzanine_path, done_file_name) Path(done_file_path).touch() return {} @activity_task_logger([ job_io_fields.H264_MEZZANINE_OUTPUT_S3_KEY, job_io_fields.PRORES_MEZZANINE_METADATA, job_io_fields.UPC, ], is_last_handler=False, is_first_handler=True) def transfer_h264_from_s3_to_ripper(inputs): """Transfer H.264 file from s3 to ripper windows share. Args: inputs (dict): Inputs. Returns: dict: empty. """ return _transfer_h264_from_s3_to_ripper_storage(inputs) @activity_task_logger([ job_io_fields.H264_MEZZANINE_OUTPUT_S3_KEY, job_io_fields.PRORES_MEZZANINE_METADATA, job_io_fields.UPC, ]) def transfer_h264_from_s3_to_ripper_storage(inputs): """Transfer H.264 file from s3 to ripper windows share. Args: inputs (dict): Inputs. Returns: dict: empty. """ return _transfer_h264_from_s3_to_ripper_storage(inputs) def _transfer_h264_from_s3_to_ripper_storage(inputs): """Transfer H.264 file from s3 to ripper windows share. Args: inputs (dict): Inputs. Returns: dict: empty. """ h264_mezzanine_path = ( inputs[job_io_fields.H264_MEZZANINE_OUTPUT_S3_KEY]) video_s3_bucket = assets.VIDEO_S3_BUCKET filename = get_ripper_h264_filename(inputs) on_prem_h264_mezzanine_path = os.path.join( config.RIPPER_OUTPUT_VIDEO_WINDOWS_MOUNT_PATH, config.OUTPUT_VIDEO_H264_MEZZANINE_DIR) try: os.mkdir(on_prem_h264_mezzanine_path) except FileExistsError: pass windows_storage_path = os.path.join( on_prem_h264_mezzanine_path, filename) s3_client = s3.get_s3_client() s3_client.download_file( video_s3_bucket, h264_mezzanine_path, windows_storage_path) return {} @activity_task_logger([ job_io_fields.UPC, job_io_fields.TIFF_THUMBNAIL_S3_KEY ], is_last_handler=False, is_first_handler=True) def transfer_tiff_thumbnail_from_s3_to_ripper(inputs): """Transfer TIFF from S3 to Ripper. Args: inputs (dict): Inputs. Returns: dict: empty. """ tiff_thumbnail_path = ( inputs[job_io_fields.TIFF_THUMBNAIL_S3_KEY]) tiff_s2_thumbnail_s3_key = ( assets.TIFF_S2_THUMBNAIL_FILENAME_S3_KEY_TEMPLATE.format( tiff_thumbnail_path=tiff_thumbnail_path)) tiff_s5_thumbnail_s3_key = ( assets.TIFF_S5_THUMBNAIL_FILENAME_S3_KEY_TEMPLATE.format( tiff_thumbnail_path=tiff_thumbnail_path)) upc = inputs[job_io_fields.UPC] s2_filename = assets.TIFF_S2_THUMBNAIL_FILENAME_RIPPER_TEMPLATE.format( upc=upc) s5_filename = assets.TIFF_S5_THUMBNAIL_FILENAME_RIPPER_TEMPLATE.format( upc=upc) on_prem_thumbnail_path = os.path.join( config.RIPPER_OUTPUT_THUMBNAIL_WINDOWS_MOUNT_PATH, config.OUTPUT_TIFF_DIR) try: os.mkdir(on_prem_thumbnail_path) except FileExistsError: pass windows_s2_storage_path = os.path.join(on_prem_thumbnail_path, s2_filename) windows_s5_storage_path = os.path.join(on_prem_thumbnail_path, s5_filename) s3_client = s3.get_s3_client() s3_client.download_file( assets.VIDEO_S3_BUCKET, tiff_s2_thumbnail_s3_key, windows_s2_storage_path) s3_client.download_file( assets.VIDEO_S3_BUCKET, tiff_s5_thumbnail_s3_key, windows_s5_storage_path) return {} @activity_task_logger( [ job_io_fields.UPC, job_io_fields.TRANSFER_PRORES_FROM_S3_TO_RIPPER_JOB_ID ], is_first_handler=False, is_last_handler=False, job_type=transfer_prores_from_s3_to_ripper.__name__ ) def is_ripper_mezzanine_processing_done(inputs): """Check if ripper mezzanine process is done. Args: inputs (dict): Inputs. Returns: dict: Outputs. """ return _ripper_asset_processing_check( inputs, job_io_fields.TRANSFER_PRORES_FROM_S3_TO_RIPPER_JOB_ID, ASSET_TYPE_PRORES, job_io_fields.IS_RIPPER_MEZZANINE_PROCESSING_DONE ) @activity_task_logger( [ job_io_fields.UPC, job_io_fields.TRANSFER_H264_FROM_S3_TO_RIPPER_JOB_ID ], is_first_handler=False, is_last_handler=False, job_type=transfer_h264_from_s3_to_ripper.__name__ ) def is_ripper_h264_processing_done(inputs): """Check if ripper H.264 process is done. Args: inputs (dict): Inputs. Returns: dict: Outputs. """ return _ripper_asset_processing_check( inputs, job_io_fields.TRANSFER_H264_FROM_S3_TO_RIPPER_JOB_ID, ASSET_TYPE_H264, job_io_fields.IS_RIPPER_H264_PROCESSING_DONE ) @activity_task_logger( [ job_io_fields.UPC, job_io_fields.TRANSFER_TIFF_THUMBNAIL_FROM_S3_TO_RIPPER_JOB_ID ], is_first_handler=False, is_last_handler=False, job_type=transfer_tiff_thumbnail_from_s3_to_ripper.__name__ ) def is_ripper_tiff_thumbnail_processing_done(inputs): """Check if ripper TIFF thumbnails process is done. Args: inputs (dict): Inputs. Returns: dict: Outputs. """ return _ripper_asset_processing_check( inputs, job_io_fields.TRANSFER_TIFF_THUMBNAIL_FROM_S3_TO_RIPPER_JOB_ID, [ASSET_TYPE_TIFF_S2, ASSET_TYPE_TIFF_S5], job_io_fields.IS_RIPPER_TIFF_THUMBNAIL_PROCESSING_DONE ) def _ripper_asset_processing_check(inputs, job_key, asset_types, result_key): """Check if asset was processed.""" job_id = inputs[job_key] upc = inputs[job_io_fields.UPC] if isinstance(asset_types, list): _asset_types = asset_types else: _asset_types = [asset_types] is_ready = all( ows_video.is_asset_ready_for_delivery(upc, asset_type) for asset_type in _asset_types ) if is_ready: ows_video.change_job_status(job_id, job_statuses.COMPLETE) return {result_key: is_ready}