#!/usr/bin/env python3 """ Hook Transcode Test Pipeline - Modular Version Orchestrates video processing through multiple Modal workers and MediaConvert """ import json import logging from pathlib import Path import click from pipeline.constants import ( AWS_REGION, DATABASE_URL_PROD, DATABASE_URL_STAGING, DEST_BUCKET, GLOCKENSPIEL_REPO, RAW_UPLOADS_BUCKET, SUNO_DATA_UPLOADS_BUCKET, ) from pipeline.database import DatabaseClient from pipeline.hook_video_worker import HookVideoWorker from pipeline.mediaconvert import MediaConvertTester from pipeline.s3_operations import S3Client from pipeline.video_upload_worker import VideoUploadWorker # Set up logging logger = logging.getLogger(__name__) logging.basicConfig( level=logging.INFO, format="%(asctime)s - %(name)s - %(levelname)s - %(message)s" ) # Create outputs directory OUTPUT_DIR = Path("outputs") OUTPUT_DIR.mkdir(exist_ok=True) def get_database_url(db_env: str) -> str: """Get the database URL for the specified environment""" if db_env.lower() == "prod": if not DATABASE_URL_PROD: raise ValueError("DATABASE_URL_PROD not configured in environment") return DATABASE_URL_PROD elif db_env.lower() == "staging": if not DATABASE_URL_STAGING: raise ValueError("DATABASE_URL_STAGING not configured in environment") return DATABASE_URL_STAGING else: raise ValueError( f"Invalid database environment: {db_env}. Use 'prod' or 'staging'" ) def construct_raw_upload_key(hook_data): """Construct S3 key from the VideoUpload data""" upload_id = hook_data.get("upload_id") or hook_data.get("raw_video_upload_id") original_ext = hook_data.get("original_file_ext", "mp4") if upload_id and original_ext: return f"raw_uploads/{upload_id}.{original_ext}" else: raise ValueError("No valid S3 key could be constructed") def run_pipeline( hook_id: str, db_env: str = "staging", skip_upload: bool = False, skip_hook_gen: bool = False, skip_mediaconvert: bool = False, ): """Run the complete pipeline based on hook data Args: hook_id: The UUID of the hook to test db_env: Database environment to use ('prod' or 'staging') skip_upload: Skip the video upload worker step skip_hook_gen: Skip the hook video generation step skip_mediaconvert: Skip the MediaConvert step """ logger.info("Starting Hook Transcode Pipeline") logger.info(f"Hook ID: {hook_id}") logger.info(f"Database: {db_env.upper()}") logger.info(f"S3 Bucket: {SUNO_DATA_UPLOADS_BUCKET}") # Get deployment type from environment or use default import os deployment_type = os.getenv("DEPLOYMENT_TYPE", "dev") logger.info(f"Deployment: {deployment_type}") results = {} # Create hook-specific output directory hook_output_dir = OUTPUT_DIR / hook_id hook_output_dir.mkdir(parents=True, exist_ok=True) logger.info(f"Output directory: {hook_output_dir}") # Initialize clients db_client = DatabaseClient() s3_client = S3Client(AWS_REGION) upload_worker = VideoUploadWorker(deployment_type) hook_worker = HookVideoWorker(deployment_type) mediaconvert_tester = MediaConvertTester(s3_client, GLOCKENSPIEL_REPO) try: # Step 1: Connect to database and get hook data logger.info("=" * 50) logger.info("Step 1: Fetch Hook Data") database_url = get_database_url(db_env) db_client.connect(database_url, db_env) hook_data = db_client.get_hook_data(hook_id) results["hook_data"] = { "hook_id": hook_data.get("hook_id"), "raw_video_upload_id": hook_data.get("raw_video_upload_id"), "upload_id": hook_data.get("upload_id"), "video_s3_id": hook_data.get("video_s3_id"), "original_file_ext": hook_data.get("original_file_ext"), } # Construct S3 key raw_upload_key = construct_raw_upload_key(hook_data) s3_client.download_file( s3_bucket=RAW_UPLOADS_BUCKET, s3_key=raw_upload_key, local_path=hook_output_dir / f"0_raw_upload.mp4", replace_existing=False, ) logger.info(f"Using source video: {raw_upload_key}") results["source_s3_key"] = raw_upload_key # Step 3: Run video upload worker upload_output_key = None if not skip_upload: logger.info("=" * 50) logger.info("Step 3: Video Upload Worker") # Always use the original S3 key from the hook data for video upload worker # since it's already in S3 upload_result = upload_worker.process_video(raw_upload_key, hook_data) results["upload_worker"] = upload_result # Download processed video if available upload_output_key = upload_result.get("output_s3_key") if upload_output_key: try: local_upload = hook_output_dir / f"1_upload_processed.mp4" s3_client.download_file( s3_bucket=SUNO_DATA_UPLOADS_BUCKET, s3_key=upload_output_key, local_path=local_upload, ) results["downloaded_upload"] = str(local_upload) except Exception as e: logger.warning(f"Could not download upload output: {e}") else: logger.info("=" * 50) logger.info("Step 3: Video Upload Worker - SKIPPED") logger.info("Using raw upload video directly") # Step 4: Run hook video generation worker hook_output_key = None if not skip_hook_gen: logger.info("=" * 50) logger.info("Step 4: Hook Video Generation") hook_result = hook_worker.generate_video( upload_output_key or raw_upload_key, hook_data ) results["hook_gen_worker"] = hook_result # Download generated hook video hook_output_key = hook_result.get("output_s3_key") if hook_output_key: try: local_hook = hook_output_dir / f"2_hook_generated.mp4" s3_client.download_file( s3_bucket=SUNO_DATA_UPLOADS_BUCKET, s3_key=hook_output_key, local_path=local_hook, ) results["downloaded_hook"] = str(local_hook) except Exception as e: logger.warning(f"Could not download hook output: {e}") else: logger.info("=" * 50) logger.info("Step 4: Hook Video Generation - SKIPPED") logger.info("Using upload processed video directly") # Step 5: Run MediaConvert test if not skip_mediaconvert: logger.info("=" * 50) logger.info("Step 5: MediaConvert Test") mediaconvert_result = mediaconvert_tester.test_video( hook_id=hook_id, input_s3_key=hook_output_key or upload_output_key or raw_upload_key, skip_copy=skip_hook_gen, ) results["mediaconvert"] = mediaconvert_result # Download MediaConvert output (HLS manifest) if successful if mediaconvert_result.get("status") == "COMPLETE": try: # Download and concatenate 720p version logger.info("Downloading 720p HLS segments...") concat_result = mediaconvert_tester.download_and_concat_720p( DEST_BUCKET, mediaconvert_result["output_prefix"], hook_output_dir, ) if concat_result: results["downloaded_720p_mp4"] = concat_result logger.info( f"Downloaded and concatenated 720p MP4: {concat_result}" ) except Exception as e: logger.warning( f"Could not download HLS manifest (may not exist yet): {e}" ) else: logger.info("=" * 50) logger.info("Step 5: MediaConvert Test - SKIPPED") # Final summary logger.info("=" * 50) logger.info("Pipeline Complete") logger.info("Pipeline Results:") logger.info(json.dumps(results, indent=2)) logger.info(f"Downloaded files saved to: {hook_output_dir}") return results except Exception as e: logger.error(f"Pipeline failed: {e}") logger.error("Partial Results:") logger.error(json.dumps(results, indent=2)) raise finally: db_client.close() @click.command() @click.argument("hook_id") @click.option("--deployment", default="dev", help="Deployment type (dev/staging/prod)") @click.option( "--db", default="staging", type=click.Choice(["staging", "prod"], case_sensitive=False), help="Database environment to use", ) @click.option("--skip-upload", is_flag=True, help="Skip the video upload worker step") @click.option( "--skip-hook-gen", is_flag=True, help="Skip the hook video generation step" ) @click.option("--skip-mediaconvert", is_flag=True, help="Skip the MediaConvert step") def main( hook_id: str, deployment: str, db: str, skip_upload: bool, skip_hook_gen: bool, skip_mediaconvert: bool, ): """ Run the Hook Transcode Test Pipeline on an existing hook from the database. HOOK_ID: The UUID of the video hook to test """ # Set deployment type in environment for this run import os os.environ["DEPLOYMENT_TYPE"] = deployment run_pipeline(hook_id, db, skip_upload, skip_hook_gen, skip_mediaconvert) if __name__ == "__main__": main()