""" PYTHONPATH=. env/bin/python3.6 dev/backfill_s3_object_id.py -bucket=qa-orcd-podcast-output-assets -asset_transcoder_api_url=https://qa-ows-asset-transcoder.theorchard.io for prod: use https://workstation.theorchard.com/grass/asset-transcoder and an OA grass token """ import argparse import multiprocessing import os import requests import traceback import sys import boto3 from botocore.exceptions import ClientError from podcast.models.episode import Episode from podcast.connectors import mysql request_headers = { 'Grass-Account-Type': 'vendor', 'Grass-Account-Id': '1', 'Correlation-Id': '12', 'Orchard-Identity-UUID': 'podcast-admin', 'Content-Type': 'application/json', # 'session': 'e39ee7bf6b7c6745be5ea3a28fb6be9f' } s3_client = boto3.client('s3') resource = boto3.resource('s3') def update_file_metadata(response, key, bucket, metadata): """Update a file's metadata using the REPLACE directive. Args: key (str): filename with the extension bucket (str): bucket metadata (dict): key / value pairs to be merged with the file's existing metadata """ asset_metadata = { 'ContentType': response['ContentType'], 'Metadata': response['Metadata'], 'MetadataDirective': 'REPLACE' } asset_metadata.update(metadata) if 'Metadata' in metadata: asset_metadata['Metadata'].update(response['Metadata']) copy_source = { 'Bucket': bucket, 'Key': key } resource.meta.client.copy( copy_source, bucket, key, ExtraArgs=asset_metadata, SourceClient=s3_client ) def get_episode_ids(): """Return all the 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). Returns: dict: containing the paginated episodes. """ with mysql.pod_db_session(read_only=True) as session: query = session.query(Episode.id).filter(Episode.is_deleted.isnot(True)) rows = query.all() return [row.id for row in rows] def get_asset(episode_id, url): assets = requests.get( '{}/assets-by-ids-and-types'.format(url), params={ 'object_ids': episode_id, 'object_types': 'episode' }, headers=request_headers ).json() for asset in assets['items']: if asset['media_type'] == 'audio': return '{}.{}'.format(asset['filename'], asset['asset_type'].lower()) def update_episode_asset_metadata(arguments): episode_id, bucket, asset_transcoder_api_url = arguments asset = get_asset(episode_id, asset_transcoder_api_url) if not asset: return episode_id, True try: s3_input_file = s3_client.head_object(Bucket=bucket, Key=asset) metadata = s3_input_file['Metadata'] if not ('object_type' in metadata and 'object_id' in metadata): print(episode_id) update_file_metadata(s3_input_file, asset, bucket, { 'Metadata': { 'object_id': str(episode_id), 'object_type': 'episode' } }) return episode_id, True except ClientError as e: error_message = 'error updating metadata for {key} in bucket {bucket}'.format(key=asset, bucket=bucket) print(error_message) return episode_id, False def do_it(asset_transcoder_api_url, bucket): episode_ids = get_episode_ids() print(len(episode_ids)) pool = multiprocessing.Pool(12) worker_input = list( zip( episode_ids, [bucket] * len(episode_ids), [asset_transcoder_api_url] * len(episode_ids) ) ) iterator = pool.imap_unordered(update_episode_asset_metadata, worker_input) while True: try: episode_id, done = next(iterator) except multiprocessing.TimeoutError: continue except StopIteration: break except Exception: print('Failed metadata {}'.format(episode_id)) # Print traceback because we can't reraise it here traceback.print_exc(file=sys.stdout) else: print('Finished metadata {} {}'.format(episode_id, done)) pool.close() pool.join() def main(): """Extract shell arguments and start ingest.""" argparser = argparse.ArgumentParser(prog='thing') argparser.add_argument( '-asset_transcoder_api_url', required=False, help='asset_transcoder_api_url, https://qa-ows-asset-transcoder.theorchard.io' ) argparser.add_argument( '-bucket', required=False, help='bucket, dev-orcd-asset-transcoder-input' ) args = argparser.parse_args() asset_transcoder_api_url = args.asset_transcoder_api_url or 'http://localhost:5001' bucket = args.bucket or 'dev-orcd-asset-transcoder-input' do_it(asset_transcoder_api_url, bucket) if __name__ == '__main__': main()