"""Logic Tier for Episode.""" from concurrent.futures import ThreadPoolExecutor import copy from datetime import datetime from functools import partial from flask import g from podcast import config from podcast.connectors import mysql from podcast.connectors.sentry import send_to_sentry from podcast.constants import asset_types as asset_types_consts from podcast.constants import common as common_constants from podcast.constants import error from podcast.constants.episode_replication_status import AD_LOCATIONS_REPLICATION_COMPLETED from podcast.constants.episode_replication_status import AD_LOCATIONS_REPLICATION_FAILED from podcast.constants.episode_replication_status import ASSET_REPLICATION_COMPLETED from podcast.constants.episode_replication_status import ASSET_REPLICATION_FAILED from podcast.constants.episode_replication_status import EPISODE_CREATED from podcast.constants.episode_replication_status import EPISODE_CREATION_FAILED from podcast.constants.episode_replication_status import REPLICATION_TYPE_SINGLE from podcast.constants.feature_flag import FEATURE_PODCAST_IA_RESTRUCTURE from podcast.constants.feature_flag import FEATURE_PODCAST_SERIAL_EPISODIC from podcast.logic import email from podcast.logic import megaphone from podcast.logic import transcription as transcription_logic from podcast.logic import user as user_logic from podcast.logic import user_v2 as user_v2_logic from podcast.models import episode as episode_model from podcast.models import episode_replication_status as episode_replication_status_model from podcast.models import insertion_point from podcast.models import ows_asset_transcoder as oat from podcast.models import podcast as podcast_model from podcast.models import s3 from podcast.utils import api_utils from podcast.utils import feature_flag_utils from podcast.utils import signed_urls from podcast.utils.exc import OwsError def create_episode(podcast_id, data): """Create episode function. Args: podcast_id (int): The podcast unique identifier Returns: dict: containing a dict with the created episode. """ user_logic.current_user_has_read_only_access_then_raise() podcast = _get_podcast_or_exception(podcast_id) if feature_flag_utils.get_feature_flag(FEATURE_PODCAST_SERIAL_EPISODIC): data = _validate_season_and_episode_number_or_exception(podcast, data) data['trailer_type'] = _validate_trailer_type(data) created_episode = episode_model.create_episode({'podcast_id': podcast_id}) return update_episode(podcast_id, created_episode['id'], data, podcast=podcast) def update_episode(podcast_id, episode_id, data, podcast=None): """Update episode function. Args: episode_id (int): The unique identifier of the episode podcast_id (int): The unique identifier of the podcast data (dict): The payload with episode metadata podcast (dict): optional dict of podcast data Returns: dict: containing a dict with the updated episode. """ user_logic.current_user_has_read_only_access_then_raise() if podcast is None: # episode being updated rather than created podcast = _get_podcast_or_exception(podcast_id) if feature_flag_utils.get_feature_flag(FEATURE_PODCAST_SERIAL_EPISODIC): data = _validate_season_and_episode_number_or_exception(podcast, data) data['trailer_type'] = _validate_trailer_type(data) if 'artwork_filename' in data: oat.commit_or_delete_asset(data.pop('artwork_filename'), episode_id, 'artwork', 'episode') if 'audio_filename' in data: insertion_point.delete_all(episode_id) transcription_logic.delete_by_episode_id(episode_id) audio_filename = data.pop('audio_filename') oat.commit_or_delete_asset(audio_filename, episode_id, 'audio', 'episode') old_episode = episode_model.get_episode_by_id(episode_id) updated_episode = episode_model.update_episode(episode_id, data) if data['draft']: assets = _draft_assets(episode_id) else: assets = _published_assets(episode_id) if _did_scheduled_date_change(old_episode, updated_episode): email.send_episode_scheduled(updated_episode) return create_or_update_megaphone( podcast, updated_episode, assets ) def reupload_mp3_to_megaphone(episode_id): """Reupload mp3 to megaphone by organization admin. Args: episode_id (int): The unique identifier of the episode Returns: dict: containing a dict with the updated episode. """ user_logic.current_user_is_org_admin_or_raise() episode = episode_model.get_episode_by_id(episode_id) episode_megaphone_id = episode['megaphone_id'] if not episode_megaphone_id: raise OwsError.bad_request(error.ERROR_MESSAGE_MEGAPHONE_EPISODE_NOT_FOUND) if episode['draft']: audio_url = _draft_assets(episode_id).get('audio_url') else: audio_url = _published_assets(episode_id).get('audio_url') if not audio_url: raise OwsError.bad_request(error.ERROR_ASSET_FINAL_NOT_FOUND) podcast = podcast_model.get_podcast_by_id(episode['podcast_id']) megaphone.reupload_episode_mp3( podcast['megaphone_id'], episode_megaphone_id, podcast['network_id'], audio_url, episode_id ) return {'episode_id': episode_id} def get_episodes(podcast_id, limit=0, offset=0, filter_by_state='ALL', order_by='published_date', sort_order='asc', start_date=None, end_date=None): """Get all episodes for a podcast. Args: podcast_id (int): The podcast unique identifier limit (int): how many podcasts to retrieve. offset (int): the offset (for pagination). filter_by_state (str): filtering options order_by (str): order by field name sort_order (str): sorting order asc|desc start_date (str): publish date after start date or None end_date (str): publish date before end date or None Returns: dict: paginated dictionary """ podcast = podcast_model.get_podcast_by_id(podcast_id) if feature_flag_utils.get_feature_flag(FEATURE_PODCAST_IA_RESTRUCTURE): user_v2_logic.current_user_owns_show_family_or_raise(podcast['show_family_id']) else: user_logic.current_user_owns_podcast_or_raise(podcast) episodes = episode_model.get_episodes( podcast_id, limit, offset, filter_by_state, order_by, sort_order, start_date, end_date) total_drafts = episode_model.get_drafts_count(podcast_id) total_done = episode_model.get_done_count(podcast_id) return { 'items': episodes['items'], 'pagination': { 'total_records': total_done + total_drafts, 'total_done': total_done, 'total_drafts': total_drafts, } } def get_episode_by_id(podcast_id, episode_id): """Get episode by it's podcast_id and episode_id. Args: podcast_id (int): The unique identifier of podcast episode_id (int): The unique identifier of episode Returns: dict: containing a dict with episode data """ podcast = podcast_model.get_podcast_by_id(podcast_id) if feature_flag_utils.get_feature_flag(FEATURE_PODCAST_IA_RESTRUCTURE): user_v2_logic.current_user_owns_show_family_or_raise(podcast['show_family_id']) else: user_logic.current_user_owns_podcast_or_raise(podcast) return episode_model.get_episode_by_id(episode_id) def delete_episode_from_megaphone(podcast_id, episode_megaphone_id): """Delete episode from megaphone. Args: podcast_id (int): The unique identifier of podcast episode_megaphone_id (str): The unique identifier of episode in megaphone """ podcast = podcast_model.get_podcast_by_id(podcast_id) podcast_megaphone_id = podcast['megaphone_id'] network_id = podcast['network_id'] if not podcast_megaphone_id: raise OwsError.bad_request(error.ERROR_MESSAGE_INVALID_PODCAST_ID) megaphone.delete_episode(podcast_megaphone_id, episode_megaphone_id, network_id) def delete_episode(podcast_id, episode_id): """Delete episode by it's podcast_id and episode_id. Args: podcast_id (int): The unique identifier of podcast episode_id (int): The unique identifier of episode Returns: dict: containing a dict with episode data """ user_logic.current_user_is_org_admin_or_raise() episode_payload = episode_model.get_episode_by_id(episode_id) with mysql.pod_db_session() as session: episode_result = episode_model.delete_episode(podcast_id, episode_id, session) episode_megaphone_id = episode_payload['megaphone_id'] if episode_megaphone_id: delete_episode_from_megaphone(podcast_id, episode_megaphone_id) return episode_result def _published_assets(episode_id): """Check episode assets to be valid.""" assets = oat.get_episode_assets(episode_id, asset_types_consts.MAX_AUDIO_SIGNED_URL_DURATION) if ('audio_path' not in assets) or ('audio_duration' not in assets): raise OwsError.bad_request(error.ERROR_ASSET_FINAL_NOT_FOUND) s3.check_s3_file_exists( bucket_name=config.OUTPUT_ASSETS_BUCKET_NAME, file_key=assets['audio_path'] ) if asset_types_consts.SUBTYPE_XLARGE_COVER in assets: s3.check_s3_file_exists( bucket_name=config.OUTPUT_ASSETS_BUCKET_NAME, file_key=assets[asset_types_consts.SUBTYPE_XLARGE_COVER] ) return assets def _draft_assets(episode_id): """Get draft assets.""" try: assets = oat.get_episode_assets(episode_id, asset_types_consts.MAX_AUDIO_SIGNED_URL_DURATION) except OwsError as err: if err.status != 404: raise assets = [] return assets def create_planned_inventory(podcast_id, data): """Create planned inventory and send to megaphone. Args: podcast_id (int): The podcast unique identifier data (dict): The dates and expected ads count Returns: list(dict): containing a dict with the podcast id """ user_logic.current_user_has_read_only_access_then_raise() podcast = podcast_model.get_podcast_by_id(podcast_id) if feature_flag_utils.get_feature_flag(FEATURE_PODCAST_IA_RESTRUCTURE): user_v2_logic.current_user_owns_show_family_or_raise(podcast['show_family_id']) else: user_logic.current_user_owns_podcast_or_raise(podcast) episodes = episode_model.create_planned_inventory(podcast_id, data) for episode_data in episodes['items']: create_or_update_megaphone(podcast, episode_data, []) return episodes def create_or_update_megaphone(podcast_data, episode_data, assets): """Create or update in megaphone, then update the episode with returned info.""" if not episode_data['is_reviewed'] or podcast_data['feed_type'] in \ [common_constants.APPLE_SUBSCRIPTION, common_constants.YOUTUBE]: return episode_data megaphone_podcast_id = podcast_data['megaphone_id'] network_id = podcast_data['network_id'] megaphone_response = megaphone.create_or_update_episode(megaphone_podcast_id, episode_data, assets, network_id) if not episode_data['megaphone_id']: return episode_model.update_episode(episode_data['id'], { 'megaphone_id': megaphone_response['id'], 'megaphone_uid': megaphone_response['uid'] }) return episode_data def get_episode_is_processing(podcast_id, episode_id): """Get if episode audio is processing in megaphone.""" podcast = podcast_model.get_podcast_by_id(podcast_id) if feature_flag_utils.get_feature_flag(FEATURE_PODCAST_IA_RESTRUCTURE): user_v2_logic.current_user_owns_show_family_or_raise(podcast['show_family_id']) else: user_logic.current_user_owns_podcast_or_raise(podcast) episode = episode_model.get_episode_by_id(episode_id) processing = megaphone.is_megaphone_audio_processing( podcast['megaphone_id'], episode['megaphone_id'], podcast['network_id'] ) return {'is_processing': processing} def get_episodes_by_ids(ids): """Get episodes by ids. Args: ids (list): The episode unique identifiers. Returns: list(dict): containing a list of dicts with episodes metadata. """ episodes = episode_model.get_episodes_by_ids(ids) podcast_ids = [episode['podcast_id'] for episode in episodes['items']] podcasts = podcast_model.get_podcasts_by_ids(podcast_ids) if feature_flag_utils.get_feature_flag(FEATURE_PODCAST_IA_RESTRUCTURE): user_v2_logic.current_user_owns_show_families_or_raise(podcasts['items']) else: user_logic.current_user_owns_podcasts_or_raise(podcasts['items']) return episodes def get_episode_wav_asset(episode_id): """Get an episode's signed audio wav url from output bucket. Only admins are permitted. Get episode audio wav asset upload whose status is enoding_completed from oat. Format audio_filename with filename and asset type from OAT response, will be used as key in s3 buckets. Check if the audio wav file already exists in the output bucket. If not, copy audio wav file from the input to output bucket. Generate signed url from output_bucket_cdn and audio_filename using cloudfront_signer having validity for 10 mins. Args: episode_id (int): The episode unique identifier. Returns: dict: contains audio_wav_url. """ user_logic.current_user_is_org_admin_or_raise() wav_asset = oat.get_object_asset_by_id_and_type(episode_id, 'episode', asset_types_consts.TYPE_FILE_WAV) audio_filename = '{}.{}'.format(wav_asset['filename'], wav_asset['asset_type'].lower()) try: s3.check_s3_file_exists(config.OUTPUT_ASSETS_BUCKET_NAME, audio_filename) except OwsError: s3.copy_s3_file( config.INPUT_ASSETS_BUCKET_NAME, audio_filename, config.OUTPUT_ASSETS_BUCKET_NAME, audio_filename ) audio_wav_url = signed_urls.sign_url(api_utils.asset_url(audio_filename)) return {'audio_wav_url': audio_wav_url} def get_most_recent(limit): """Get most recent published episodes for a user. Args: limit (int): max num to return Returns: list(int): list of episode ids """ if feature_flag_utils.get_feature_flag(FEATURE_PODCAST_IA_RESTRUCTURE): podcast_ids = user_v2_logic.podcast_ids_owned_by_current_user() else: podcast_ids = user_logic.podcast_ids_owned_by_current_user() return episode_model.get_most_recent_episodes_by_podcast_ids(podcast_ids, limit) def _get_podcast_or_exception(podcast_id): podcast = podcast_model.get_podcast_by_id(podcast_id) if feature_flag_utils.get_feature_flag(FEATURE_PODCAST_IA_RESTRUCTURE): user_v2_logic.current_user_owns_show_family_or_raise(podcast['show_family_id']) else: user_logic.current_user_owns_podcast_or_raise(podcast) return podcast def _validate_season_and_episode_number_or_exception(podcast, data): episode_type = data.get('episode_type') if data['draft'] is False and \ episode_type in (common_constants.EPISODE_TYPE_BONUS, common_constants.EPISODE_TYPE_FULL): if data.get('season_number') is None: raise OwsError.bad_request(error.ERROR_MISSING_SEASON_NUMBER) if data.get('episode_number') is None: raise OwsError.bad_request(error.ERROR_MISSING_EPISODE_NUMBER) clean_data = data.copy() if episode_type == common_constants.EPISODE_TYPE_TRAILER: clean_data['season_number'] = None clean_data['episode_number'] = None return clean_data def _did_scheduled_date_change(old_episode, new_episode): if new_episode['draft']: return False if old_episode['published_date'] == new_episode['published_date']: return False if new_episode['published_date'] < datetime.utcnow(): return False return True def _validate_trailer_type(data): """Validate trailer type when episode type is trailer.""" trailer_type = None if data.get('episode_type') == common_constants.EPISODE_TYPE_TRAILER: if data['draft'] is False and data.get('trailer_type') is None: raise OwsError.bad_request(error.ERROR_MISSING_TRAILER_TYPE) trailer_type = data.get('trailer_type') return trailer_type def _replicate_episode_for_type_bulk( original_podcast_id, podcast_id, seasons, replication_type, user_id, is_copy_apple_episode_id, episode_data): """Replicate episode and create entries in episode_replication_status_model respectively.""" season_id = None original_podcast_id = original_podcast_id podcast_id = podcast_id original_episode_id = episode_data.get('id') original_season_number = episode_data.get('season_number') published_date = episode_data.get('published_date') # Fetches season_id of the matched episode's season_number with new season's number if seasons and original_season_number: season_id = next((season['id'] for season in seasons if season['number'] == original_season_number), None) response_object = { 'original_podcast_id': original_podcast_id, 'podcast_id': podcast_id, 'original_episode_id': original_episode_id, 'episode_id': None } episode_payload = { 'podcast_id': podcast_id, 'title': episode_data.get('title'), 'description': episode_data.get('description'), 'episode_type': episode_data.get('episode_type'), 'trailer_type': episode_data.get('trailer_type'), 'content': episode_data.get('content'), 'season_number': episode_data.get('season_number'), 'episode_number': episode_data.get('episode_number'), 'external_id': episode_data.get('external_id'), 'published_date': published_date.strftime('%Y-%m-%d %H:%M:%S') if published_date else None, 'is_reviewed': False, 'season_id': season_id, 'apple_id': episode_data.get('apple_id') if is_copy_apple_episode_id else None } with mysql.pod_db_session() as session: try: episode_replication_status_response = episode_replication_status_model.create_status({ 'original_podcast_id': original_podcast_id, 'podcast_id': podcast_id, 'original_episode_id': original_episode_id, 'replication_type': replication_type }, session) created_episode = episode_model.create_episode(episode_payload, session, user_id) response_object['episode_id'] = created_episode['id'] episode_replication_status_model.update_status({ 'id': episode_replication_status_response['id'], 'episode_id': created_episode['id'], 'status': EPISODE_CREATED }, session) except Exception as e: if common_constants.EPISODE_REPLICATION_STATUS_RESPONSE not in locals(): response_object['error_message'] = error.ERROR_MESSAGE_CANNOT_RECORD_STATUS.format(e.status, e.message) elif common_constants.CREATED_EPISODE not in locals() and \ common_constants.EPISODE_REPLICATION_STATUS_RESPONSE in locals(): episode_replication_status_model.update_status({ 'id': episode_replication_status_response['id'], 'status': EPISODE_CREATION_FAILED }) response_object['error_message'] = error.ERROR_MESSAGE_CANNOT_CREATE_EPISODE.format(e.status, e.message) return response_object def _replicate_ad_locations_for_type_bulk(assets, user_id, processed_episode): """Replicate ad locations. Only replicate ad locations for newly created episode whose asset creation is successful. """ replicated_asset = next( ( asset for asset in assets if asset['episode_id'] == processed_episode['episode_id'] and asset['status'] == ASSET_REPLICATION_COMPLETED ), None ) if replicated_asset: try: insertion_points = insertion_point.get_insertion_points(processed_episode['original_episode_id'])['items'] if insertion_points: new_insertion_points = [ { 'episode_id': processed_episode['episode_id'], 'point_type': insertion_point['point_type'], 'timecode': insertion_point['timecode'], 'count': insertion_point['count'] } for insertion_point in insertion_points ] insertion_point.create_insertion_points( processed_episode['episode_id'], new_insertion_points, user_id=user_id) replicated_asset['status'] = AD_LOCATIONS_REPLICATION_COMPLETED episode_replication_status_model.update_status(replicated_asset) except Exception as err: replicated_asset['status'] = AD_LOCATIONS_REPLICATION_FAILED episode_replication_status_model.update_status(replicated_asset) send_to_sentry(err, {}, 500, error.ERROR_MESSAGE_FAILED_EPISODE_AD_LOCATIONS) processed_episode['status'] = replicated_asset['status'] return processed_episode def replicate_episodes_and_create_assets( original_podcast_id, podcast_id, episode_ids, seasons, replication_type, is_copy_apple_episode_id, is_copy_ad_locations ): """Make copies of episodes and replicate assets.""" processed_episodes = [] try: episodes = get_episodes_by_ids(episode_ids) episode_list = episodes['items'] user_id = g.current_user['id'] replicate_episodes_func = partial( _replicate_episode_for_type_bulk, original_podcast_id, podcast_id, seasons, replication_type, user_id, is_copy_apple_episode_id) with ThreadPoolExecutor(max_workers=common_constants.MAX_WORKERS) as executor: processed_episodes = list(executor.map(replicate_episodes_func, episode_list)) processed_episode_ids = [ {'original_episode_id': episode['original_episode_id'], 'episode_id': episode['episode_id']} for episode in processed_episodes if episode['episode_id'] ] assets_replication_response = oat.replicate_episodes_audio_assets(processed_episode_ids) if assets_replication_response.status_code > 399: errored_episodes = [ {**episode_obj, 'status': ASSET_REPLICATION_FAILED} for episode_obj in processed_episode_ids] with ThreadPoolExecutor(max_workers=common_constants.MAX_WORKERS) as executor: executor.map(episode_replication_status_model.update_status, errored_episodes) raise OwsError( status=assets_replication_response.status_code, message=assets_replication_response.text) assets = assets_replication_response.json()['items'] with ThreadPoolExecutor(max_workers=common_constants.MAX_WORKERS) as executor: executor.map(episode_replication_status_model.update_status, assets) processed_episodes_with_ads = [] if is_copy_ad_locations: for episode in processed_episodes: processed_episodes_with_ads.append(_replicate_ad_locations_for_type_bulk(assets, user_id, episode)) return processed_episodes_with_ads except Exception as e: send_to_sentry(e.message, error.ERROR_MESSAGE_FAILED_EPISODE_OR_ASSETS_REPLICATION, e.status, e) return processed_episodes def _replicate_episode_for_type_single( original_podcast, original_episode, podcasts, user_id, podcast_id): """Replicate single episode and create entries in episode_replication_status_model respectively.""" with mysql.pod_db_session() as session: try: original_podcast_id = original_podcast['id'] original_episode_id = original_episode['id'] destination_podcast = next((podcast for podcast in podcasts if podcast['id'] == podcast_id), None) processed_episode = { 'original_podcast_id': original_podcast_id, 'original_episode_id': original_episode_id, 'episode_id': None } published_date = original_episode.get('published_date') if not destination_podcast: raise OwsError.bad_request(error.ERROR_MESSAGE_INVALID_PODCAST_ID) is_copy_apple_episode_id = original_podcast['feed_type'] == common_constants.PUBLIC_RSS and \ destination_podcast['feed_type'] == common_constants.APPLE_SUBSCRIPTION destination_podcast_id = destination_podcast['id'] processed_episode['podcast_id'] = destination_podcast_id processed_episode['feed_type'] = destination_podcast['feed_type'] processed_episode['original_podcast_feed_type'] = original_podcast['feed_type'] episode_payload = { 'podcast_id': destination_podcast_id, 'title': original_episode['title'], 'description': original_episode['description'], 'episode_type': original_episode['episode_type'], 'trailer_type': original_episode.get('trailer_type'), 'content': original_episode['content'], 'season_number': original_episode['season_number'], 'episode_number': original_episode['episode_number'], 'external_id': original_episode.get('external_id'), 'published_date': published_date.strftime('%Y-%m-%d %H:%M:%S') if published_date else None, 'apple_id': original_episode.get('apple_id') if is_copy_apple_episode_id else None, 'is_reviewed': False } if original_podcast['show_type'] == common_constants.SHOW_TYPE_SERIAL: seasons = destination_podcast['seasons'] season_id = next( (season['id'] for season in seasons if season['number'] == original_episode['season_number']), None ) episode_payload['season_id'] = season_id if destination_podcast['feed_type'] == common_constants.APPLE_SUBSCRIPTION and \ original_podcast['feed_type'] == common_constants.PUBLIC_RSS: episode_payload['apple_id'] = original_episode.get('apple_id') episode_replication_status_response = episode_replication_status_model.create_status( { 'original_podcast_id': original_podcast_id, 'podcast_id': destination_podcast_id, 'original_episode_id': original_episode_id, 'replication_type': REPLICATION_TYPE_SINGLE }, session) created_episode = episode_model.create_episode(episode_payload, session, user_id) processed_episode['episode_id'] = created_episode['id'] processed_episode['status'] = EPISODE_CREATED episode_replication_status_model.update_status( { 'id': episode_replication_status_response['id'], 'episode_id': created_episode['id'], 'status': EPISODE_CREATED }, session) except Exception as e: if 'destination_podcast_id' not in locals(): processed_episode['error_message'] = f'{e.message}, podcast_id: {podcast_id}' elif common_constants.EPISODE_REPLICATION_STATUS_RESPONSE not in locals(): processed_episode['error_message'] = error.ERROR_MESSAGE_CANNOT_RECORD_STATUS.format( e.status, e.message) elif common_constants.CREATED_EPISODE not in locals() and \ common_constants.EPISODE_REPLICATION_STATUS_RESPONSE in locals(): episode_replication_status_model.update_status( { 'id': episode_replication_status_response['id'], 'status': EPISODE_CREATION_FAILED }) processed_episode['error_message'] = error.ERROR_MESSAGE_CANNOT_CREATE_EPISODE.format( e.status, e.message) return processed_episode def _replicate_ad_locations(original_episode_id, copy_ad_locations_podcast_ids, assets, user_id, processed_episodes): """Replicate ad locations. Replicate ad locations only for requested podcast ids whose feed-type is public-rss and destination podcast is also of feed-type public-rss and replicated episode's asset replication is successful. """ replicated_assets = [] processed_public_rss_episodes = [ episode for episode in processed_episodes if episode['feed_type'] == episode['original_podcast_feed_type'] == common_constants.PUBLIC_RSS ] for processed_public_rss_episode in processed_public_rss_episodes: if processed_public_rss_episode['podcast_id'] in copy_ad_locations_podcast_ids: replicated_assets = [ asset for asset in assets if asset['status'] == ASSET_REPLICATION_COMPLETED and processed_public_rss_episode.get( 'episode_id') == asset['episode_id'] ] if replicated_assets: try: insertion_points = insertion_point.get_insertion_points(original_episode_id)['items'] if insertion_points: with mysql.pod_db_session() as session: for replicated_asset in replicated_assets: replicated_asset_episode_id = replicated_asset['episode_id'] new_insertion_points = [ { 'episode_id': replicated_asset_episode_id, 'point_type': insertion_point['point_type'], 'timecode': insertion_point['timecode'], 'count': insertion_point['count'], 'created_by': user_id } for insertion_point in insertion_points ] insertion_point.create_insertion_points( replicated_asset_episode_id, new_insertion_points, session) replicated_asset['status'] = AD_LOCATIONS_REPLICATION_COMPLETED episode_replication_status_model.update_status(replicated_asset) except Exception: for replicated_asset in replicated_assets: replicated_asset['status'] = AD_LOCATIONS_REPLICATION_FAILED episode_replication_status_model.update_status(replicated_asset) for processed_episode in processed_episodes: processed_episode['status'] = next( ( replicated_asset['status'] for replicated_asset in replicated_assets if replicated_asset['episode_id'] == processed_episode.get('episode_id') ), processed_episode['status'] ) return processed_episodes def replicate_single_episode_and_assets(data): """Replicate single episode and it's asset to multiple podcasts.""" user_logic.current_user_has_read_only_access_then_raise() original_episode_id = data['original_episode_id'] original_episode = episode_model.get_episode_by_id(original_episode_id) original_podcast_id = original_episode['podcast_id'] podcast_ids = copy.deepcopy(data['podcast_ids']) podcast_ids.append(original_podcast_id) podcasts = podcast_model.get_podcasts_by_ids(podcast_ids)['items'] user_v2_logic.current_user_owns_show_families_or_raise(podcasts) processed_episodes = [] user_id = g.current_user['id'] copy_ad_locations_podcast_ids = data.get('copy_ad_locations_podcast_ids', None) original_podcast = next((podcast for podcast in podcasts if podcast['id'] == original_podcast_id)) try: # episode replication replicate_episode_func = partial( _replicate_episode_for_type_single, original_podcast, original_episode, podcasts, user_id) with ThreadPoolExecutor(max_workers=common_constants.MAX_WORKERS) as executor: processed_episodes = list(executor.map(replicate_episode_func, data['podcast_ids'])) processed_episode_ids = [ {'original_episode_id': episode['original_episode_id'], 'episode_id': episode['episode_id']} for episode in processed_episodes if episode['episode_id'] ] # asset replication oat api post call assets_replication_response = oat.replicate_episodes_audio_assets(processed_episode_ids) # handles asset replication oat api failure if assets_replication_response.status_code > 399: for processed_episode in processed_episodes: if processed_episode.get('episode_id'): processed_episode['status'] = ASSET_REPLICATION_FAILED errored_episodes = [ {**episode_obj, 'status': ASSET_REPLICATION_FAILED} for episode_obj in processed_episode_ids] with ThreadPoolExecutor(max_workers=common_constants.MAX_WORKERS) as executor: executor.map(episode_replication_status_model.update_status, errored_episodes) raise OwsError( status=assets_replication_response.status_code, message=assets_replication_response.text) assets = assets_replication_response.json()['items'] # updates status in episode replication status table and api response with ThreadPoolExecutor(max_workers=common_constants.MAX_WORKERS) as executor: executor.map(episode_replication_status_model.update_status, assets) for processed_episode in processed_episodes: is_matching_episode = next( (asset for asset in assets if asset.get('episode_id') == processed_episode.get('episode_id')), None ) if is_matching_episode: processed_episode['status'] = is_matching_episode['status'] # ad location replication if copy_ad_locations_podcast_ids: processed_episodes = _replicate_ad_locations( original_episode_id, copy_ad_locations_podcast_ids, assets, user_id, processed_episodes) except Exception as e: send_to_sentry(e.message, error.ERROR_MESSAGE_FAILED_EPISODE_OR_ASSETS_REPLICATION, e.status, e) return processed_episodes