"""Copy Podcasting Assets. Run this in dev with command: PYTHONPATH=. env/bin/python3.11 scripts/copy_podcast_assets.py -env dev -team sme-affiliate-podcasts -copy_artwork yes """ import argparse import boto3 import csv import logging import os import requests import sys import time from botocore.exceptions import ClientError from podcast.utils import api_utils from podcast.models import episode as episode_model from podcast.models import podcast as podcast_model from podcast.models import network as network_model from podcast.api import app from flask import g from podcast.utils.exc import OwsError # --- Logging setup --- log = logging.getLogger('main') log.setLevel(logging.DEBUG) fmt = logging.Formatter('%(asctime)s - %(levelname)s - %(message)s') sh = logging.StreamHandler(sys.stdout) sh.setFormatter(fmt) log.addHandler(sh) # --- Global Config --- users = {'dev': '366', 'qa': '1298', 'prod': '88'} request_headers = None s3 = boto3.client('s3') PodcastIdMismatchError = type("PODCASTIDMISMATCH", (Exception,), {}) # --- Asset & S3 Utilities --- def get_object_asset_by_id_and_type(object_id, object_type, asset_type): """"Get Object asset by id and type. Args: object_id : object_id object_type : object_type asset_type : object_type """ response = requests.get( f'{oat_api_url}/asset-by-id-and-type', params={'object_id': object_id, 'object_type': object_type, 'asset_type': asset_type}, headers=request_headers ) if response.status_code > 399: raise OwsError.response_error(response) return response.json() def _get_assets(object_id, object_type, signed_url_duration=10): response = requests.get( f'{oat_api_url}/assets-by-ids-and-types', params={ 'object_ids': [object_id], 'object_types': [object_type], 'signed_url_duration': signed_url_duration }, headers=request_headers ) if response.status_code > 399: raise OwsError.response_error(response) return response.json()['items'] def get_episode_mp3_file(episode_id): """"Get episode mp3 filename. Args: episode_id : episode_id """ assets = _get_assets(episode_id, 'episode') for asset in assets: for final in asset['asset_finals']: if final['asset_type'] == 'MP3': return final['filename'] return None def get_episode_artwork_filename(episode_id): """"Get episode artwork filename. Args: episode_id : episode_id """ assets = _get_assets(episode_id, 'episode') return next( ( f"{item['filename']}.{item['asset_type'].lower()}" for item in assets if item.get('asset_type') in ['JPG', 'TIF'] and 'filename' in item ), None ) def object_exists(bucket, key, max_retries=5, base_delay=1): """Check if an object exists in S3 with exponential backoff on failure. Args: bucket (str): S3 bucket name. key (str): S3 object key. max_retries (int): Maximum number of retries on failure. base_delay (float): Initial delay between retries (doubles on each retry). Returns: bool: True if object exists, False if 404, raises otherwise. """ attempt = 0 while attempt <= max_retries: try: s3.head_object(Bucket=bucket, Key=key) return True except ClientError as e: error_code = e.response["Error"]["Code"] if error_code == "404": return False elif error_code == "400": wait = base_delay * (2 ** attempt) time.sleep(wait) attempt += 1 continue else: raise except Exception: raise raise RuntimeError(f"Max retries exceeded for HEAD s3://{bucket}/{key}") def copy_s3_object(source_bucket, source_key, dest_bucket, dest_key, asset_label): """"Copy s3 object. Args: source_bucket : bucket, source_key : source_key, dest_bucket : dest_bucket, asset_label : asset_label """ try: log.info(f"Copying {asset_label}: s3://{source_bucket}/{source_key} -> s3://{dest_bucket}/{dest_key}") s3.copy_object( Bucket=dest_bucket, CopySource={'Bucket': source_bucket, 'Key': source_key}, Key=dest_key ) return True except Exception as e: log.warning(f"[MISSING] {asset_label} - s3://{source_bucket}/{source_key}") log.error(f"Error : {e}") return False # --- Main Copy Logic --- def copy_podcast_assets(filename, team, copy_artwork): """"Copy podcast assets. Args: filename : filename, team : team, copy_artwork : copy_artwork """ if not os.path.exists(filename): raise FileNotFoundError(f'File not found: {filename}') log.info(f'Reading file: {filename}\n') podcasts = podcast_model.get_all_podcasts().get('items') networks = network_model.get_networks().get('items') error_episode_count = 0 errored_episode_ids = [] copy_failures = { 'podcast_artwork': [], 'episode_artwork': [], 'episode_wav': [], 'episode_mp3': [] } with open(filename) as csv_file: reader = csv.DictReader(csv_file) with app.app_context(): with app.test_request_context('/dummy'): for row_number, row in enumerate(reader, start=2): # Start from 2 to account for header try: episode_id = int(row['Episode Episode ID']) podcast_id = int(row['Podcast Podcast ID']) print("\n-----------------------------------------------------------------------------") log.info(f"\nRow {row_number} | Podcast ID: {podcast_id} | Episode ID: {episode_id}") episode = episode_model.get_episode_by_id(episode_id) episode_title = episode.get('title') episode_podcast_id = episode.get('podcast_id') episode_season_number= episode.get('season_number') episode_number = episode.get('episode_number') if episode_podcast_id != podcast_id: raise PodcastIdMismatchError( ( f"Podcast ID mismatch: Episode {episode_id} -> " f"Expected {episode_podcast_id}, Got {podcast_id}" ) ) podcast_title, network_id = next( ((p['title'], p['network_id']) for p in podcasts if p['id'] == episode_podcast_id), (None, None) ) network_name = next((n['name'] for n in networks if n['id'] == network_id), None) episode_wav_asset = get_object_asset_by_id_and_type(episode_id, 'episode', 'WAV') episode_wav = f"{episode_wav_asset['filename']}.{episode_wav_asset['asset_type'].lower()}" episode_mp3 = get_episode_mp3_file(episode_id) base_path = f"{team}/{network_name}/{podcast_title}" if episode_season_number: prefix = f"Season{episode_season_number}_Episode{episode_number}" else: prefix = f"Episode{episode_number}" episode_path = f"{base_path}/{prefix}_{episode_title}" asset_filename = f"{podcast_title}_{prefix}_{episode_title}" if copy_artwork == 'yes': try: podcast_artwork_asset = get_object_asset_by_id_and_type( podcast_id, 'podcast', ['TIF', 'JPG'] ) podcast_artwork = ( f"{podcast_artwork_asset['filename']}." f"{podcast_artwork_asset['asset_type'].lower()}" ) podcast_artwork_filename = ( f"{podcast_title}.{podcast_artwork_asset['asset_type'].lower()}" ) if not copy_s3_object( s3_buckets['raw_assets_bucket'], podcast_artwork, s3_buckets['vendor_bucket'], f"{base_path}/{podcast_artwork_filename}", "Podcast Artwork" ): copy_failures['podcast_artwork'].append({ 'row': row_number, 'podcast_id': podcast_id, 'key': podcast_artwork }) except Exception as e: log.warning(f"Failed to copy podcast artwork: {e}") episode_artwork = get_episode_artwork_filename(episode_id) if episode_artwork: artwork_ext = episode_artwork.split('.')[-1] if not copy_s3_object( s3_buckets['raw_assets_bucket'], episode_artwork, s3_buckets['vendor_bucket'], f"{episode_path}/{asset_filename}.{artwork_ext}", "Episode Artwork" ): copy_failures['episode_artwork'].append({ 'row': row_number, 'episode_id': episode_id, 'key': episode_artwork }) if episode_wav and not copy_s3_object( s3_buckets['raw_assets_bucket'], episode_wav, s3_buckets['vendor_bucket'], f"{episode_path}/{asset_filename}.wav", "Episode WAV" ): copy_failures['episode_wav'].append({ 'row': row_number, 'episode_id': episode_id, 'key': episode_wav }) if episode_mp3 and not copy_s3_object( s3_buckets['final_assets_bucket'], episode_mp3, s3_buckets['vendor_bucket'], f"{episode_path}/{asset_filename}.mp3", "Episode MP3" ): copy_failures['episode_mp3'].append({ 'row': row_number, 'episode_id': episode_id, 'key': episode_mp3 }) log.info(f"!! Done processing Episode {episode_id}: {episode_title}") print("\n-----------------------------------------------------------------------------") except PodcastIdMismatchError as e: error_episode_count += 1 errored_episode_ids.append(episode_id) log.error(f"[Row {row_number}] {e}") except Exception as e: error_episode_count += 1 errored_episode_ids.append(episode_id) log.exception(f"[Row {row_number}] Unexpected error for Episode {episode_id}: {e}") # Summary print(f"\nErrored episodes: {error_episode_count} -> {errored_episode_ids}") for key, failures in copy_failures.items(): print(f"Missing {key.replace('_', ' ')}: {failures}") # --- Entry Point --- def main(): """Script entry point. Executes the main function to copy podcast assets. """ global request_headers global s3_buckets global oat_api_url parser = argparse.ArgumentParser(prog='Copy Podcast Assets') parser.add_argument('-env', required=True, help='Environment (dev, qa, prod)') parser.add_argument('-team', required=True, help='Team name (used in S3 path)') parser.add_argument('-copy_artwork', required=True, help='yes/no to copy artwork') args = parser.parse_args() env = args.env team = args.team copy_artwork = args.copy_artwork.lower() filename = os.environ.get('FILENAME', 'copyPodcastAssets.csv') oat_api_url = f'https://{env}-ows-asset-transcoder.theorchard.io' api_utils.get_user_id = lambda: users[env] request_headers = { 'Grass-Account-Type': 'vendor', 'Grass-Account-Id': '1', 'Correlation-Id': '12', 'Orchard-Identity-UUID': 'podcast-admin', 'Content-Type': 'application/json' } s3_buckets = { 'raw_assets_bucket': f'{env}-orcd-asset-transcoder-input', 'final_assets_bucket': f'{env}-orcd-podcast-output-assets', 'vendor_bucket': f'{env}-podcasting-assets' } log.info(f"Starting asset copy in '{env}' for team '{team}' (copy artwork: {copy_artwork})") with app.app_context(): g.log = log copy_podcast_assets(filename, team, copy_artwork) if __name__ == '__main__': main()