"""Megaphone Show Importer. run this with docker compose run --build megaphone-show-importer --env=dev --podcast_id={podcast_id} --network_id={network_id} --api_key={api_key} --episode_assets={skip/ingest} """ import argparse from collections import Counter from datetime import timezone import ntpath import os import sys import time from urllib.parse import urlparse import boto3 from dateutil import parser import ffmpy from PIL import Image import requests from podcast import config from podcast.connectors import mysql from podcast.constants.common import SHOW_TYPE_SERIAL from podcast.models import category as category_model from podcast.models import episode as episode_model from podcast.models import insertion_point as insertion_point_model from podcast.models import podcast as podcast_model from podcast.models import podcast_links as podcast_links_model from podcast.models import podcast_season as season_model from podcast.models.api_episode import ApiEpisode from podcast.models.api_podcast import ApiPodcast from podcast.models.show_family import create_show_family from podcast.utils import api_utils from podcast.utils.exc import OwsError users = { 'dev': '366', 'qa': '1298', 'prod': '88' } ASSETS_DIR = './dev/rss_ingester_assets/' IMAGE = 'jpg' AUDIO = 'wav' request_headers = None def _get_full_filename(token_response_body, asset_type): return '{}.{}'.format( token_response_body.get('filename'), asset_type ) def poll(filename, asset_transcoder_api_url): """Poll for asset finish processing api call.""" global request_headers status_response = requests.get( '{}/status/{}'.format(asset_transcoder_api_url, filename), headers=request_headers ).json() if 'status' in status_response: if status_response['status'].endswith('_error'): print(status_response) print('for filename {}'.format(filename)) sys.exit(100) if status_response['status'] == 'encoding_completed': return status_response['assets'] time.sleep(4) return poll(filename, asset_transcoder_api_url) def create_asset(url, object_id, object_type, asset_type, asset_transcoder_api_url): """Create asset api call.""" global request_headers if not url: return try: asset_response = requests.get(url) if asset_response.status_code > 299: return asset_data = asset_response.content except requests.exceptions.RequestException: return ext = os.path.splitext(ntpath.basename(urlparse(url).path))[1] begin_filename = '{}{}{}'.format(ASSETS_DIR, object_type, ext) with open(begin_filename, 'wb') as handler: handler.write(asset_data) if asset_type == IMAGE: image = Image.open(begin_filename) image = image.resize((3000, 3000)) image = image.convert('RGB') end_filename = '{}{}.{}'.format(ASSETS_DIR, object_type, 'jpg') content_type = 'image/jpeg' image.save(end_filename) else: end_filename = '{}{}.{}'.format(ASSETS_DIR, object_type, 'wav') content_type = 'audio/wav' try: ff = ffmpy.FFmpeg( inputs={begin_filename: None}, outputs={end_filename: '-acodec pcm_s16le -ar 44100 -ac 2 -y'} ) ff.run() except ffmpy.FFRuntimeError: return original_filename = ntpath.basename(urlparse(url).path) original_filename = os.path.splitext(original_filename)[0] + '.' + asset_type try: token_response = requests.get('{}/upload-token'.format(asset_transcoder_api_url), headers=request_headers) if token_response.status_code > 299: return token_response_body = token_response.json() except requests.exceptions.RequestException: return upload_creds = token_response_body.get('credentials') bucket_name = token_response_body.get('bucket') dest_file_name = _get_full_filename(token_response_body, asset_type) file_to_upload = './{}'.format(end_filename) client = boto3.client( 's3', aws_access_key_id=upload_creds.get('aws_access_key_id'), aws_secret_access_key=upload_creds.get('aws_secret_access_key'), aws_session_token=upload_creds.get('token') ) s3_metadata = { 'object_type': object_type, 'object_id': str(object_id), 'asset_type': asset_type.upper(), 'original_filename': original_filename, } client.upload_file( file_to_upload, bucket_name, dest_file_name, ExtraArgs={ 'Metadata': s3_metadata, 'ContentType': content_type } ) return dest_file_name def create_new_podcast_links(podcast_id, guid): """Create website embed on new podcast.""" store_id = podcast_links_model.get_website_embed_store()['id'] url = 'https://playlist.megaphone.fm?p={}'.format(guid) link = ''.format(url) podcast_links_model.create_podcast_link( {'podcast_id': podcast_id, 'store_id': store_id, 'link': link}) def create_orchard_podcast(podcast, network_id, asset_transcoder_api_url, env): """Orchestrate podcast api calls from the parsed megaphone data.""" explicit = 'clean' if podcast['explicit']: explicit = 'explicit' categories = category_model.get_categories()['items'] podcast_categories = [category['id'] for category in categories if category['name'] in podcast['itunesCategories']] podcast_title = podcast['title'] if env == 'qa': podcast_title = '[PROD] {}'.format(podcast['title']) orch_podcast_dict = { 'network_id': network_id, 'title': podcast_title, 'description': podcast['summary'], 'slug': podcast['slug'], 'host': podcast['author'], 'owner': podcast['ownerName'], 'copyright': podcast['copyright'], 'link': podcast['link'], 'email': podcast['ownerEmail'], 'explicit': explicit, 'show_type': podcast['podcastType'], 'created_date': parser.parse(podcast['createdAt']), 'updated_date': parser.parse(podcast['updatedAt']), 'is_deleted': False, 'language': podcast['language'], 'megaphone_id': podcast['id'], 'categories': podcast_categories, } orchard_show_family = create_show_family({ 'title': podcast_title, 'network_id': network_id, 'created_by': users[env], 'updated_by': users[env] }) orch_podcast_dict['show_family_id'] = orchard_show_family['id'] orchard_podcast = podcast_model.create_podcast(orch_podcast_dict) create_asset(podcast['imageFile'], orchard_podcast['id'], 'podcast', IMAGE, asset_transcoder_api_url) create_new_podcast_links(orchard_podcast['id'], podcast['uid']) return orchard_podcast def create_orchard_insertion_points(episode, orch_episode_id, duration): """Orchestrate insertion point api calls from the parsed megaphone data.""" insertion_data = [] if episode['preCount'] is None: episode['preCount'] = 0 if episode['postCount'] is None: episode['postCount'] = 0 if episode['preCount'] > 0: insertion_data.append({ 'point_type': 'pre', 'timecode': episode['preOffset'], 'count': episode['preCount'] }) if episode['postCount'] > 0: insertion_data.append({ 'point_type': 'post', 'timecode': float(duration) - float(episode['postOffset']), 'count': episode['postCount'] }) if len(episode['insertionPoints']) > 0: counter = Counter(episode['insertionPoints']) for timecode in counter.keys(): insertion_data.append({ 'point_type': 'mid', 'timecode': timecode, 'count': counter[timecode] }) if len(insertion_data) > 0: insertion_point_model.create_insertion_points(orch_episode_id, insertion_data) def _get_episode_data(episode, episode_description_filter_string): explicit = 'clean' if episode['explicit']: explicit = 'explicit' planned_pre_roll_count = episode['expectedAdhash'].count('0') planned_mid_roll_count = episode['expectedAdhash'].count('1') planned_post_roll_count = episode['expectedAdhash'].count('2') if episode['draft']: draft = True else: draft = False pub_date = None if episode['pubdate']: pub_date = parser.parse(episode['pubdate']) pub_date = pub_date.astimezone(timezone.utc) pub_date = pub_date.strftime('%Y-%m-%dT%H:%M:%S') episode_info = { 'title': episode['title'], 'description': episode['summary'], 'season_id': episode.get('seasonId'), 'season_number': episode['seasonNumber'], 'episode_number': episode['episodeNumber'], 'episode_type': episode['episodeType'], 'content': explicit, 'planned_pre_roll_count': planned_pre_roll_count, 'planned_mid_roll_count': planned_mid_roll_count, 'planned_post_roll_count': planned_post_roll_count, 'created_date': parser.parse(episode['createdAt']).astimezone(timezone.utc), 'updated_date': parser.parse(episode['updatedAt']).astimezone(timezone.utc), 'is_deleted': False, 'megaphone_id': episode['id'], 'draft': draft, 'published_date': pub_date, 'megaphone_uid': episode['uid'] } if episode_description_filter_string: episode_info['description'] = episode['summary'].replace(episode_description_filter_string, '') if episode['guid']: # if the episode has a guid use that for the uuid episode_info['uuid'] = episode['guid'] return episode_info def _update_episode_assets(episode, orch_episode_id, asset_transcoder_api_url): create_asset(episode['imageFile'], orch_episode_id, 'episode', IMAGE, asset_transcoder_api_url) filename = create_asset(episode['audioFile'], orch_episode_id, 'episode', AUDIO, asset_transcoder_api_url) duration = None if filename: assets = poll(filename, asset_transcoder_api_url) duration = [asset['duration'] for asset in assets if asset['asset_type'] == 'FLAC'][0] / 1000 create_orchard_insertion_points(episode, orch_episode_id, duration or episode['duration']) def _get_seasons(podcast_id): """Get seasons by podcast_id.""" with mysql.pod_db_session(read_only=True) as session: rows = session.query(season_model.PodcastSeason).filter( season_model.PodcastSeason.podcast_id == podcast_id ).all() return [row.to_dict() for row in rows] def get_or_create_orch_seasons(podcast_id, season_number): """Get or create new seasons. Get existing seasons by podcast_id. Check if season exists by season number from megaphone episode. If exists return the matching season. Else, check for the max existing season number. In case when seasons doesn't exists between max existing season number and the current season number in megaphone episode data, create all the missing seasons and the season for season number in megaphone episode. Season name is given by default. Return all seasons which includes newly created seasons. """ existing_seasons = _get_seasons(podcast_id) matching_season = [season for season in existing_seasons if season['number'] == season_number] if matching_season: return matching_season else: max_existing_season_number = max((season['number'] for season in existing_seasons), default=0) while max_existing_season_number < season_number: existing_seasons.append({ 'number': max_existing_season_number + 1, 'name': f'Season {max_existing_season_number + 1}' }) max_existing_season_number += 1 orch_seasons = season_model.create_seasons(podcast_id, existing_seasons, True) return orch_seasons def create_orchard_episode( orch_podcast, episode, asset_transcoder_api_url, with_episode_assets, episode_description_filter_string): """Create orchard episode from the parsed megaphone data.""" orch_data = _get_episode_data(episode, episode_description_filter_string) orch_data['podcast_id'] = orch_podcast['id'] orch_episode = episode_model.create_episode(orch_data) if with_episode_assets: _update_episode_assets(episode, orch_episode['id'], asset_transcoder_api_url) return orch_episode def update_orchard_episode( episode, orch_episode, asset_transcoder_api_url, with_assets, episode_description_filter_string): """Update orchard episode from the parsed megaphone data.""" orch_data = _get_episode_data(episode, episode_description_filter_string) orch_episode = episode_model.update_episode(orch_episode['id'], orch_data) if with_assets: _update_episode_assets(episode, orch_episode['id'], asset_transcoder_api_url) return orch_episode def main(): """Extract shell arguments and start ingest.""" global request_headers argparser = argparse.ArgumentParser(prog='megaphone tester') argparser.add_argument('--api_key', required=True, help='API Key') argparser.add_argument('--episode_assets', required=True, help='skip or ingest') argparser.add_argument('--episode_id', required=False, help='episode id') argparser.add_argument('--podcast_id', required=True, help='podcast id') argparser.add_argument('--network_id', required=True, help='network id') argparser.add_argument('--env', required=True, help='env', default='qa') argparser.add_argument( '--asset_transcoder_api_url', required=False, help='asset_transcoder_api_url, https://qa-ows-asset-transcoder.theorchard.io' ) argparser.add_argument( '--episode_description_filter_string', required=False, help='string to be removed from the episode title' ) args = argparser.parse_args() asset_transcoder_api_url = args.asset_transcoder_api_url or 'http://localhost:5001' network_id = args.network_id podcast_id = args.podcast_id episode_id = args.episode_id config.MEGAPHONE_API_TOKEN = args.api_key with_episode_assets = True if args.episode_assets == 'ingest' else False env = args.env episode_description_filter_string = args.episode_description_filter_string def get_user_id(): """Mock get user id return value.""" return users[env] api_utils.get_user_id = get_user_id request_headers = { 'Grass-Account-Type': 'vendor', 'Grass-Account-Id': '1', 'Correlation-Id': '12', 'Orchard-Identity-UUID': 'podcast-admin', 'Content-Type': 'application/json' } try: orch_podcast = podcast_model.get_podcast_by_megaphone_id(podcast_id) orch_episodes = episode_model.get_episodes(orch_podcast['id']).get('items') orch_episodes_ids = [ep.get('megaphone_id') for ep in orch_episodes] except OwsError: megaphone_api = ApiPodcast(network_id) podcast = megaphone_api.get(podcast_id) orch_podcast = create_orchard_podcast(podcast, network_id, asset_transcoder_api_url, env) orch_episodes_ids = [] orch_episodes = [] megaphone_api_episode = ApiEpisode(network_id) if episode_id: episodes = [megaphone_api_episode.get(podcast_id, episode_id)] else: page = 1 limit = 500 episodes = [] while limit: episodes_data = megaphone_api_episode.get_all( podcast_id, page, per_page=500) episodes += episodes_data['items'] total_records = episodes_data['pagination']['total_records'] if limit > total_records: limit = 0 else: limit += 500 page += 1 for episode in episodes: print('\nRunning for episode: {}.'.format(episode['title'])) if orch_podcast['show_type'] == SHOW_TYPE_SERIAL: if episode.get('seasonNumber'): orch_seasons = get_or_create_orch_seasons(orch_podcast['id'], episode['seasonNumber']) episode['seasonId'] = orch_seasons[-1]['id'] else: print('\nSeason number not found for serial show episode.\n') if episode['id'] not in orch_episodes_ids: create_orchard_episode( orch_podcast, episode, asset_transcoder_api_url, with_episode_assets, episode_description_filter_string) else: orch_episode = next((e for e in orch_episodes if e['megaphone_id'] == episode['id']), None) # only update episode that have been published since the last import newly_published = orch_episode['status'] == 'draft' and not episode['draft'] forced_episode_import = orch_episode['megaphone_id'] == episode_id with_assets = (newly_published or forced_episode_import) and with_episode_assets update_orchard_episode( episode, orch_episode, asset_transcoder_api_url, with_assets, episode_description_filter_string) if __name__ == '__main__': main()