#!/usr/bin/env python3 # /// script # requires-python = ">=3.12" # dependencies = [ # "boto3", # "psycopg2", # ] # /// import argparse import json import os import boto3 import psycopg2 import psycopg2.extras def parse_args(): p = argparse.ArgumentParser(description="Copy bots_generatedclip to DynamoDB clip-meta-heavy") p.add_argument( "--db-url", default=os.getenv("DATABASE_URL", "postgres://suno@localhost:5432/suno_studio"), help="Postgres connection URL", ) p.add_argument( "--table", default=os.getenv("DDB_TABLE", "clip-meta-heavy"), help="DynamoDB table name" ) p.add_argument( "--ddb-endpoint", default=os.getenv("DDB_ENDPOINT", "http://localhost:8123"), help="DynamoDB endpoint URL (local)", ) p.add_argument( "--region", default=os.getenv("AWS_DEFAULT_REGION", "us-east-2"), help="AWS region (for boto3 session)", ) p.add_argument("--batch-size", type=int, default=500, help="Fetch size per batch") p.add_argument("--limit", type=int, default=None, help="Limit rows copied") return p.parse_args() def to_item(row): # row fields: id (uuid), prompt_text (text), metadata (json/jsonb) clip_id = str(row["id"]) if row["id"] is not None else None prompt_text = row.get("prompt_text") metadata = row.get("metadata") prompt_value = None if isinstance(metadata, dict): prompt_value = metadata.get("prompt") elif isinstance(metadata, str): try: md = json.loads(metadata) prompt_value = md.get("prompt") if isinstance(md, dict) else None except Exception: prompt_value = None # Build Dynamo item (omit None or empty values) item = {"clipId": clip_id} if prompt_text: item["prompt_text"] = prompt_text if prompt_value: item["prompt"] = prompt_value return item def main(): args = parse_args() # Minimal creds for local Dynamo os.environ.setdefault("AWS_ACCESS_KEY_ID", "local") os.environ.setdefault("AWS_SECRET_ACCESS_KEY", "local") ddb = boto3.resource("dynamodb", region_name=args.region, endpoint_url=args.ddb_endpoint) table = ddb.Table(args.table) conn = psycopg2.connect(args.db_url) conn.autocommit = False # Server-side cursor to stream results cur_name = "bots_generatedclip_stream" with conn.cursor(name=cur_name, cursor_factory=psycopg2.extras.RealDictCursor) as cur: sql = "SELECT id, prompt_text, metadata FROM bots_generatedclip" if args.limit: sql += f" LIMIT {int(args.limit)}" cur.itersize = args.batch_size cur.execute(sql) total = 0 with table.batch_writer(overwrite_by_pkeys=["clipId"]) as batch: while True: rows = cur.fetchmany(args.batch_size) if not rows: break for r in rows: item = to_item(r) if not item.get("clipId"): continue batch.put_item(Item=item) total += 1 print(f"Copied {total} items to DynamoDB table '{args.table}'") conn.close() if __name__ == "__main__": main()