"""Module saves metadata about uploaded asset.""" from concurrent.futures import ThreadPoolExecutor from datetime import datetime from functools import partial import os from asset_transcoder import config from asset_transcoder.connectors import mysql from asset_transcoder.constants import asset_status as asset_status_constants from asset_transcoder.constants import asset_types as asset_types_constants from asset_transcoder.constants import asset_upload as asset_upload_constants from asset_transcoder.constants import error from asset_transcoder.constants.asset_upload import MAX_WORKERS from asset_transcoder.logic import ownership from asset_transcoder.models import asset_final as asset_final_model from asset_transcoder.models import asset_status as asset_status_model from asset_transcoder.models import asset_upload as asset_upload_model from asset_transcoder.utils import api_utils from asset_transcoder.utils import asset_type as asset_type_util from asset_transcoder.utils import s3 from asset_transcoder.utils import signed_urls from asset_transcoder.utils.exceptions import OwsError def update_asset_upload(payload): """Update asset_upload data with new values. Args: payload (dict): Asset upload data. Returns: dict: Success or error message """ payload['asset_type'] = asset_type_util.map_image_asset_type(payload.get('asset_type').upper()) clean_filename, ext = os.path.splitext(payload['filename']) payload.pop('filename') asset = asset_upload_model.get_asset_upload(clean_filename) if payload['object_id'] and payload['object_type']: _ensure_ownership_or_raise(asset.user_id, payload['object_type'], payload['object_id']) updated_asset = _update_asset_upload(asset.asset_upload_id, payload) updated_asset['config'] = api_utils.get_pipeline_config(payload['object_type']) return updated_asset def commit(filename, object_id): """Attach an asset to a particular object_id. Args: filename (string): Asset upload filename id. object_id (string): Object id to attach Returns: dict: Success or error message """ clean_filename, ext = os.path.splitext(filename) asset = asset_upload_model.get_asset_upload(clean_filename) # TODO: this ownership check is a must here. # TODO: update podcast transactional routine to handle it # TODO: https://github.com/theorchard/ows-podcast/blob/master/podcast/logic/podcast.py#L26-L40 # _ensure_ownership_or_raise(asset.user_id, asset.object_type, object_id) updated_asset = _update_asset_upload(asset.asset_upload_id, {'object_id': object_id}) s3.update_file_metadata(filename, config.INPUT_ASSETS_BUCKET_NAME, { 'Metadata': { 'object_id': str(object_id) } }) return updated_asset def _update_asset_upload(asset_id, payload): updated_asset = asset_upload_model.update_asset_upload(asset_id, payload) if payload['object_id']: asset_upload_model.delete_previous_assets( current_asset_upload_id=asset_id, object_type=updated_asset['object_type'], object_id=updated_asset['object_id'], asset_type=updated_asset['asset_type'] ) return updated_asset def delete_asset_upload(frontend_asset_type, object_id, object_type): """Delete asset_upload. Args: frontend_asset_type (str): artwork | audio object_id (str): Unigue identifier of the object in system object_type (str): podcast | episode | product | etc. Returns: dict: Success or error message """ if frontend_asset_type == 'audio': asset_types = asset_types_constants.AUDIO_FILE_TYPES else: asset_types = asset_types_constants.IMAGE_FILE_TYPES _ensure_ownership_or_raise(api_utils.get_user_id(), object_type, object_id) return asset_upload_model.delete_asset_upload(asset_types, object_id, object_type) def get_assets_by_ids_and_types(object_ids, object_types, signed_url_policy, signed_url_duration): """Get assets by ids and types.""" assets = asset_upload_model.get_objects_assets(object_ids, object_types) return { 'items': signed_urls.sign_urls(assets['items'], signed_url_policy, signed_url_duration) } def get_assets_by_ids_and_type(object_ids, object_type, signed_url_policy, signed_url_duration): """Get assets by ids and type.""" assets = asset_upload_model.get_objects_assets( object_ids, [object_type] * len(object_ids)) return { 'items': signed_urls.sign_urls(assets['items'], signed_url_policy, signed_url_duration) } def get_asset_by_id_and_type(object_id, object_type, asset_type): """Get asset whose encoding is completed by object_id, object_type and asset_type. Args: object_id (int): Unique identifier of the object in system. object_type (str): Type of object like episode or podcast. asset_type (list of str): file type like ['WAV'] for episode or ['TIF', 'JPG'] for podcast. Returns: dict: contains asset upload object. """ return asset_upload_model.get_asset_by_id_and_type(object_id, object_type, asset_type) def _replicate_episode_audio_assets_rows(assets, user_id, object_dict): """Replicate table entries for the episode asset in asset_upload, asset_status, and asset_final table. Check if asset is present based on object_id. If yes: check if final assets are present. - If yes: fetch the latest final assets i.e FLAC, DAT, MP3 rows. - Replicate entry in asset_upload table using asset and set is_duplicated to true. - Replicate entries in asset_final table using asset_finals. - Create entry in asset_status table with status as encoding_complete. - Return success dict - Else: raise final asset not found. - return error dict Else: raise asset upload not found. - Return error dict Args: assests (list of dict): list of dict containing episodes audio assets. user_id (str): uuid id of the current user. object_dict (dict): contains original_episode_id i.e object_id of existing row in asset_upload table and episode_id i.e object_id that will be linked to replicated asset_upload entry. Returns: dict: contains success or error object. """ original_episode_id = object_dict['original_episode_id'] episode_id = object_dict['episode_id'] with mysql.transcoder_db_session() as session: try: asset = next((asset for asset in assets if int(asset['object_id']) == original_episode_id), None) if asset: if asset['asset_finals']: latest_assets = {} for final_asset in asset['asset_finals']: asset_type = final_asset['asset_type'] if asset_type not in latest_assets or final_asset['id'] > latest_assets[asset_type]['id']: latest_assets[asset_type] = final_asset asset['asset_finals'] = list(latest_assets.values()) else: raise Exception(error.ERROR_ASSET_FINAL_NOT_FOUND) replicated_asset_upload = asset_upload_model.create_asset_upload( { 'user_id': user_id, 'asset_type': asset['asset_type'], 'filename': asset['filename'], 'original_filename': asset['original_filename'], 'object_id': episode_id, 'object_type': asset['object_type'], 'is_duplicated': True }, session ) for final_asset in asset['asset_finals']: asset_final_model.create_asset_final( asset_upload_id=replicated_asset_upload['id'], asset_type=final_asset['asset_type'], asset_subtype=final_asset.get('asset_subtype'), filename=final_asset['filename'], duration=final_asset.get('duration', 0), session=session ) asset_status_model.create_asset_status( asset_upload_id=replicated_asset_upload['id'], status=asset_status_constants.STATUS_ENCODING_COMPLETED, created_date=datetime.now(), errors=None, session=session ) return { 'original_episode_id': original_episode_id, 'episode_id': episode_id, 'status': asset_status_constants.STATUS_REPLICATION_COMPLETED } else: raise Exception(error.ERROR_ASSET_UPLOAD_NOT_FOUND) except Exception as err: session.rollback() return { 'original_episode_id': original_episode_id, 'episode_id': episode_id, 'status': asset_status_constants.STATUS_REPLICATION_FAILED, 'error_message': err.args[0] } def replicate_episodes_audio_assets(objects): """Replicate episodes audio assets by objects. Fetch user id from request context. Fetch assets corresponding to the object_ids, object_type=episode, and asset_type=WAV, this asset contains asset upload object and it's corresponding final assets rows i.e FLAC, DAT, MP3. Multithreading is used to execute replication using object_ids with max 10 concurrent threads. Args: objects (list of dict): contains original_episode_id and episode_id. Returns: list of dict: list item is dict that contains success or error on per entity basis. """ response_list = [] user_id = api_utils.get_user_id() original_episode_ids = [item['original_episode_id'] for item in objects] assets = asset_upload_model.get_episodes_audio_assets(original_episode_ids) replicate_func = partial(_replicate_episode_audio_assets_rows, assets, user_id) with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: for result in executor.map(replicate_func, objects): response_list.append(result) return {'items': response_list} def replicate_podcast_artwork_asset(object_dict): """Replicate podcast artwork asset. Fetch user id from request context. Fetch assets corresponding to the object_ids, object_type=podcast, and asset_type=JPG, Args: object_dict (dict): contains original_podcast_id and podcast_id. """ original_podcast_id = object_dict.get('original_podcast_id') podcast_id = object_dict.get('podcast_id') with mysql.transcoder_db_session() as session: try: original_asset_upload = asset_upload_model.get_asset_by_id_and_type( original_podcast_id, asset_upload_constants.OBJECT_TYPE_PODCAST, [asset_types_constants.TYPE_FILE_JPG]) if not original_asset_upload: raise Exception(error.ERROR_ASSET_UPLOAD_NOT_FOUND) replicated_asset_upload = asset_upload_model.create_asset_upload( { 'user_id': api_utils.get_user_id(), 'asset_type': original_asset_upload['asset_type'], 'filename': original_asset_upload['filename'], 'original_filename': original_asset_upload['original_filename'], 'object_id': podcast_id, 'object_type': original_asset_upload['object_type'], 'is_duplicated': True, }, session ) legacy_final_assets = asset_final_model.get_asset_finals(original_asset_upload['id'])['items'] for final_asset in legacy_final_assets: asset_final_model.create_asset_final( asset_upload_id=replicated_asset_upload['id'], asset_type=final_asset['asset_type'], asset_subtype=final_asset.get('asset_subtype'), filename=final_asset['filename'], duration=final_asset.get('duration', 0), session=session ) asset_status_model.create_asset_status( asset_upload_id=replicated_asset_upload['id'], status=asset_status_constants.STATUS_ENCODING_COMPLETED, created_date=datetime.now(), errors=None, session=session ) return { 'original_podcast_id': original_podcast_id, 'podcast_id': podcast_id, 'status': asset_status_constants.STATUS_REPLICATION_COMPLETED } except Exception: session.rollback() return { 'original_podcast_id': original_podcast_id, 'podcast_id': podcast_id, 'status': asset_status_constants.STATUS_REPLICATION_FAILED } def _ensure_ownership_or_raise(user_id, object_type, object_id): if not ownership.check(user_id, object_type, object_id): raise OwsError(message='User does not own object', status=401)