"""RSS ingester. Run with PYTHONPATH=. env/bin/python3.6 dev/rss_ingester.py -feed_url=https://www.briefs.fm/cocktailing.xml -slug=d -network_id=1 -timeout=30 """ import argparse import ntpath import os import sys import time from urllib.parse import urlparse import boto3 from defusedxml.ElementTree import fromstring import ffmpy from PIL import Image import requests NAMESPACE = '{http://www.itunes.com/dtds/podcast-1.0.dtd}' ASSETS_DIR = './dev/rss_ingester_assets/' request_headers = { 'Grass-Account-Type': 'vendor', 'Grass-Account-Id': '1', 'Correlation-Id': '12', 'Orchard-User-Id': 'podcast-admin', 'Content-Type': 'application/json' } IMAGE = 'jpg' AUDIO = 'wav' log_file = None def log(message): """Write to log.""" global log_file print(message) log_file.write(message + '\n') def _check_for_ffmpeg(): try: ff = ffmpy.FFmpeg( inputs={'{}a1.wav'.format(ASSETS_DIR): None}, outputs={'{}a2.wav'.format(ASSETS_DIR): '-y'} ) ff.run() except ffmpy.FFRuntimeError: pass except ffmpy.FFExecutableNotFoundError: log('FFMPEG binary not found') sys.exit(1) def get_feed_content(feed_url, timeout): """Get the feed content.""" try: response = requests.get(feed_url, timeout=timeout) except requests.exceptions.RequestException as e: log('Feed could not be found {}'.format(e)) sys.exit(1) if response.status_code > 299: log('Feed could not be found {}'.format(response.status_code)) sys.exit(1) return response.content def create_podcast_request(domain, data): """Create the podcast api call.""" response = requests.post('{}/podcasts'.format(domain), headers=request_headers, json=data) if response.status_code > 299: log('Podcast could not be created {}'.format(response.status_code)) sys.exit(1) return response.json()['id'] def create_episode_request(domain, podcast_id, data): """Create the episode api call.""" create_episode_response = requests.post( '{}/podcasts/{}/episodes'.format(domain, podcast_id), headers=request_headers, json=data ) if create_episode_response.status_code > 299: log('Episode could not be created {}'.format(create_episode_response.status_code)) return return create_episode_response.json()['id'] def get_categories_request(domain): """Get the categories api call.""" categories = requests.get('{}/categories'.format(domain), headers=request_headers) if categories.status_code > 299: log('Categories could not be fetched {}'.format(categories.status_code)) sys.exit(1) return categories.json()['items'] def _get_full_filename(token_response_body, asset_type): return '{}.{}'.format( token_response_body.get('filename'), asset_type ) def poll(filename, domain): """Poll for asset finish processing api call.""" status_response = requests.get('{}/status/{}'.format(domain, filename), headers=request_headers).json() if 'status' in status_response: if status_response['status'].endswith('_error'): log('Error transcoding asset {}'.format(filename)) return False if status_response['status'] == 'encoding_completed': return True time.sleep(4) return poll(filename, domain) def create_asset(url, object_type, asset_type, domain): """Create asset api call.""" if not url: return try: asset_response = requests.get(url) if asset_response.status_code > 299: log('Asset could not be fetched {}'.format(asset_response.status_code)) return asset_data = asset_response.content except requests.exceptions.RequestException as e: log('Feed could not be found {}'.format(e)) 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 as e: log('FFMpeg error {}'.format(e)) 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(domain), headers=request_headers) if token_response.status_code > 299: log('Could not get token {}'.format(token_response.status_code)) return token_response_body = token_response.json() except requests.exceptions.RequestException as e: log('Could not get token {}'.format(e)) 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, '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 find_category_id(orchard_categories, category_levels, category_id): """Find category id for an rss category.""" if len(category_levels) == 0: return category_id found = [x for x in orchard_categories if x['name'] == category_levels[0]] if len(found): if 'children' in found[0]: return find_category_id(found[0]['children'], category_levels[1:], found[0]['id']) else: return found[0]['id'] return category_id def _get_children_categories(category_xml, category): category.append(category_xml.get('text')) if len(category_xml) > 0: _get_children_categories(category_xml[0], category) def get_categories(channel): """Find categories from rss category.""" categories_xml = channel.findall('{}category'.format(NAMESPACE)) categories = [] for category_xml in categories_xml: category = [] _get_children_categories(category_xml, category) categories.append(category) return categories def get_podcast_info(channel): """Get the podcast info from the RSS feed.""" title = channel.find('title') link = channel.find('link') description = channel.find('description') copyright_element = channel.find('copyright') show_type = channel.find('{}type'.format(NAMESPACE)) language = channel.find('language') host = channel.find('{}author'.format(NAMESPACE)) image = channel.find('{}image'.format(NAMESPACE)) explicit = channel.find('{}explicit'.format(NAMESPACE)) owner = channel.find('{}owner'.format(NAMESPACE)) email = owner and owner.find('{}email'.format(NAMESPACE)) owner_name = owner and owner.find('{}name'.format(NAMESPACE)) return { 'title': title.text if title is not None else None, 'link': link.text if link is not None else None, 'description': description.text if description is not None else None, 'copyright': copyright_element.text if copyright_element is not None else None, 'show_type': show_type.text if show_type is not None else None, 'language': language.text if language is not None else None, 'host': host.text if host is not None else None, 'image': image.get('href') if image is not None else None, 'explicit': explicit.text if explicit is not None else None, 'categories': get_categories(channel), 'email': email.text if email is not None else None, 'owner': owner_name.text if owner_name is not None else None, } def get_episode_info(episode_xml): """Get the episode info from the RSS feed.""" title = episode_xml.find('title') description = episode_xml.find('description') season_number = episode_xml.find('{}season'.format(NAMESPACE)) episode_number = episode_xml.find('{}episode'.format(NAMESPACE)) episode_type = episode_xml.find('{}episodeType'.format(NAMESPACE)) pub_date = episode_xml.find('pubDate') explicit = episode_xml.find('{}explicit'.format(NAMESPACE)) image = episode_xml.find('{}image'.format(NAMESPACE)) audio = episode_xml.find('enclosure') return { 'title': title.text if title is not None else None, 'description': description.text if description is not None else None, 'season_number': season_number.text if season_number is not None else None, 'episode_number': episode_number.text if episode_number is not None else None, 'episode_type': episode_type.text if episode_type is not None else None, 'content': explicit.text if explicit is not None else None, 'published_date': pub_date.text if pub_date is not None else None, 'image': image.get('href') if image is not None else None, 'audio': audio.get('url') if audio is not None else None, } def create_podcast(podcast_info, network_id, slug, podcast_api_url, asset_transcoder_api_url): """Orchestrate podcast api calls from the parsed RSS data.""" categories = get_categories_request(podcast_api_url) category_ids = [] for category in podcast_info['categories']: category_id = find_category_id(categories, category, None) if category_id: category_ids.append(category_id) podcast_info['show_type'] = podcast_info['show_type'] or 'episodic' podcast_info['categories'] = category_ids podcast_info['slug'] = slug podcast_info['network_id'] = network_id if podcast_info['explicit'] == 'no': podcast_info['explicit'] = 'clean' if podcast_info['explicit'] == 'yes': podcast_info['explicit'] = 'explicit' if podcast_info['language'] == 'en-us': podcast_info['language'] = 'en' filename = create_asset(podcast_info['image'], 'podcast', IMAGE, asset_transcoder_api_url) if filename: poll(filename, asset_transcoder_api_url) podcast_info['artwork_filename'] = filename podcast_id = create_podcast_request(podcast_api_url, podcast_info) return podcast_id def create_episode(episode_info, podcast_id, podcast_api_url, asset_transcoder_api_url): """Orchestrate episode api calls from the parsed RSS data.""" episode_info['draft'] = False episode_info['episode_type'] = episode_info['episode_type'] or 'full' if episode_info['content'] == 'no': episode_info['content'] = 'clean' if episode_info['content'] == 'yes': episode_info['content'] = 'explicit' filename = create_asset(episode_info['image'], 'episode', IMAGE, asset_transcoder_api_url) if filename: poll(filename, asset_transcoder_api_url) episode_info['artwork_filename'] = filename filename = create_asset(episode_info['audio'], 'episode', AUDIO, asset_transcoder_api_url) if filename: poll(filename, asset_transcoder_api_url) episode_info['audio_filename'] = filename create_episode_request(podcast_api_url, podcast_id, episode_info) def process(feed_url, network_id, slug, podcast_api_url, asset_transcoder_api_url, timeout): """Process an entire rss feed.""" rss_string = get_feed_content(feed_url, timeout) et = fromstring(rss_string) channel = et[0] podcast_info = get_podcast_info(channel) episodes_info = [] for item in channel.findall('item'): episodes_info.append(get_episode_info(item)) podcast_id = create_podcast(podcast_info, network_id, slug, podcast_api_url, asset_transcoder_api_url) for episode_info in episodes_info: create_episode(episode_info, podcast_id, podcast_api_url, asset_transcoder_api_url) def main(): """Extract shell arguments and start ingest.""" global log_file parser = argparse.ArgumentParser(prog='megaphone tester') parser.add_argument('-feed_url', required=True, help='feed url') parser.add_argument('-slug', required=True, help='slug') parser.add_argument('-network_id', required=True, help='network id') parser.add_argument( '-podcast_api_url', required=False, help='podcast_api_url, https://qa-ows-podcast.theorchard.io' ) parser.add_argument( '-asset_transcoder_api_url', required=False, help='asset_transcoder_api_url, https://qa-ows-asset-transcoder.theorchard.io' ) parser.add_argument('-timeout', required=False, help='timeout in seconds') args = parser.parse_args() slug = args.slug feed_url = args.feed_url network_id = args.network_id timeout = args.timeout or 20 podcast_api_url = args.podcast_api_url or 'http://localhost:5000' asset_transcoder_api_url = args.asset_transcoder_api_url or 'http://localhost:5001' log_file = open('{}{}-{}-log.txt'.format(ASSETS_DIR, slug, int(time.time())), 'w') _check_for_ffmpeg() process(feed_url, network_id, slug, podcast_api_url, asset_transcoder_api_url, int(timeout)) if __name__ == '__main__': main()