""" Gather thumbnail data for video products where thumbnail selection was never completed, and output a Liquibase SQL file for a database PR to backfill thumbnail_path and thumbnail_at_milliseconds in product_video. Usage: python backfill_video_thumbnails.py [ticket_id] [sql_username] ticket_id: optional Jira ticket for the SQL changeset (default: MAINT) sql_username: optional DB username for the Liquibase changeset author (default: ) Outputs: {ticket_id}_BACKFILL_video_thumbnails.sql - Liquibase SQL to drop into database/ows_video/build/changelog/dml/ video_errors.csv - UPCs that could not be resolved with reasons Requirements: pip install boto3 pymysql python-dotenv Config: Set AR_DB_HOST, AR_DB_USER, AR_DB_PASSWORD and OWS_VIDEO_DB_HOST, OWS_VIDEO_DB_USER, OWS_VIDEO_DB_PASSWORD as env vars or in a .env file. AWS credentials must be configured (e.g. via aws-vault). """ import csv import os import re import sys from datetime import datetime, timezone import boto3 import pymysql from dotenv import load_dotenv # --- Config --- load_dotenv() ENVIRONMENT = os.environ.get("ENVIRONMENT", "qa") AWS_REGION = "us-east-1" S3_BUCKET = f"{ENVIRONMENT}-orcd-video-assets" AR_DB_HOST = os.environ.get("AR_DB_HOST", "") AR_DB_USER = os.environ.get("AR_DB_USER", "") AR_DB_PASSWORD = os.environ.get("AR_DB_PASSWORD", "") AR_DB_NAME = os.environ.get("AR_DB_NAME", "art_relations") OWS_VIDEO_DB_HOST = os.environ.get("OWS_VIDEO_DB_HOST", "") OWS_VIDEO_DB_USER = os.environ.get("OWS_VIDEO_DB_USER", "") OWS_VIDEO_DB_PASSWORD = os.environ.get("OWS_VIDEO_DB_PASSWORD", "") OWS_VIDEO_DB_NAME = os.environ.get("OWS_VIDEO_DB_NAME", "ows_video") CSV_FIELDS = [ "upc", "reason", "ingest_workflow_type", "ingest_run_datetime", "ingest_errors", "ingest_errored_jobs", "approval_errors", "approval_errored_jobs", ] def get_ar_conn(): return pymysql.connect( host=AR_DB_HOST, user=AR_DB_USER, password=AR_DB_PASSWORD, database=AR_DB_NAME, charset="utf8mb4", cursorclass=pymysql.cursors.DictCursor, ) def get_ows_video_conn(): return pymysql.connect( host=OWS_VIDEO_DB_HOST, user=OWS_VIDEO_DB_USER, password=OWS_VIDEO_DB_PASSWORD, database=OWS_VIDEO_DB_NAME, charset="utf8mb4", cursorclass=pymysql.cursors.DictCursor, ) def get_products_by_upcs(upcs): conn = get_ar_conn() placeholders = ",".join(["%s"] * len(upcs)) query = f""" SELECT release_id AS product_id, upc, latest_pipeline_run_id, latest_approval_job_id, thumbnail_path, thumbnail_at_milliseconds FROM product_video WHERE upc IN ({placeholders}) """ with conn.cursor() as cursor: cursor.execute(query, upcs) rows = cursor.fetchall() conn.close() return rows def diagnose_pipeline_run(pipeline_run_id, conn): result = {"workflow_type": None, "run_datetime": "", "job_status": None, "errors": [], "errored_jobs": []} with conn.cursor() as cursor: cursor.execute( "SELECT type, datetime FROM jobs WHERE id = %s", (pipeline_run_id,), ) row = cursor.fetchone() if row: result["workflow_type"] = row["type"] result["run_datetime"] = str(row["datetime"]) if row["datetime"] else "" cursor.execute( "SELECT status FROM job_statuses WHERE job_id = %s ORDER BY id DESC LIMIT 1", (pipeline_run_id,), ) status_row = cursor.fetchone() if status_row: result["job_status"] = status_row["status"] cursor.execute( "SELECT id FROM jobs WHERE parent_id = %s", (pipeline_run_id,), ) child_rows = cursor.fetchall() if not child_rows: return result child_ids = [r["id"] for r in child_rows] placeholders = ",".join(["%s"] * len(child_ids)) with conn.cursor() as cursor: cursor.execute( f"SELECT name, value FROM job_outputs WHERE job_id IN ({placeholders})", child_ids, ) output_rows = cursor.fetchall() cursor.execute( f""" SELECT j.type FROM jobs j JOIN job_statuses js ON js.job_id = j.id WHERE j.id IN ({placeholders}) AND js.id = (SELECT MAX(js2.id) FROM job_statuses js2 WHERE js2.job_id = j.id) AND js.status = 'ERROR' """, child_ids, ) errored_job_rows = cursor.fetchall() result["errors"] = [ f"{r['name']}: {r['value']}" if r["name"] == "error_unknown" and r["value"] else r["name"] for r in output_rows if r["name"].startswith("error_") ] result["errored_jobs"] = [r["type"] for r in errored_job_rows] return result def get_video_assets(product_id, conn): with conn.cursor() as cursor: cursor.execute( "SELECT asset_type FROM video_asset WHERE product_id = %s", (product_id,), ) return [r["asset_type"] for r in cursor.fetchall()] def parse_ms_from_path(thumbnail_path): match = re.search(r"\.(\d{7})\.jpg$", thumbnail_path) if not match: return None return int(match.group(1)) * 1000 def pick_first_frame(workflow_job_id): s3 = boto3.client("s3", region_name=AWS_REGION) prefix = f"thumbnails/original-size/{workflow_job_id}/" resp = s3.list_objects_v2(Bucket=S3_BUCKET, Prefix=prefix) contents = resp.get("Contents", []) if not contents: return None, None keys = sorted(obj["Key"] for obj in contents) # Match the frontend default: DEFAULT_THUMBNAIL = 10 (index 9), fall back to last if fewer frames chosen_key = keys[min(9, len(keys) - 1)] # thumbnail_path is stored without the leading "thumbnails/original-size/" prefix thumbnail_path = chosen_key.removeprefix("thumbnails/original-size/") thumbnail_at_ms = parse_ms_from_path(chosen_key) if thumbnail_at_ms is None: return None, None return thumbnail_path, thumbnail_at_ms def write_sql(rows, ticket_id, sql_username, output_path): with open(output_path, "w") as f: f.write("--liquibase formatted sql\n\n") f.write(f"--changeset {sql_username}:{ticket_id}\n\n") for r in rows: if r["set_path"]: f.write( f"UPDATE product_video" f" SET thumbnail_path = '{r['thumbnail_path']}'," f" thumbnail_at_milliseconds = {r['thumbnail_at_ms']}" f" WHERE release_id = {r['product_id']};\n" ) else: f.write( f"UPDATE product_video" f" SET thumbnail_at_milliseconds = {r['thumbnail_at_ms']}" f" WHERE release_id = {r['product_id']};\n" ) f.write("\n--rollback SELECT \"No Rollback\";\n") def write_skipped_csv(skipped, output_path): with open(output_path, "w", newline="") as f: writer = csv.DictWriter(f, fieldnames=CSV_FIELDS) writer.writeheader() writer.writerows(skipped) def _skip(upc, reason, ingest_workflow_type="", ingest_run_datetime="", ingest_errors="", ingest_errored_jobs="", approval_errors="", approval_errored_jobs=""): return { "upc": upc, "reason": reason, "ingest_workflow_type": ingest_workflow_type, "ingest_run_datetime": ingest_run_datetime, "ingest_errors": ingest_errors, "ingest_errored_jobs": ingest_errored_jobs, "approval_errors": approval_errors, "approval_errored_jobs": approval_errored_jobs, } def _ingest_fields(diagnosis): return { "ingest_workflow_type": diagnosis["workflow_type"] or "", "ingest_run_datetime": diagnosis["run_datetime"], "ingest_errors": ", ".join(diagnosis["errors"]), "ingest_errored_jobs": ", ".join(diagnosis["errored_jobs"]), "_ingest_job_status": diagnosis["job_status"], } def _approval_fields(diagnosis): return { "approval_errors": ", ".join(diagnosis["errors"]), "approval_errored_jobs": ", ".join(diagnosis["errored_jobs"]), "_approval_job_status": diagnosis["job_status"], } def _context_reason(errors, errored_jobs, job_status, context): if errored_jobs and errors: return f"{context} failed in {errored_jobs}: {errors}" if errored_jobs: return f"{context} failed: {errored_jobs} errored (no error output recorded)" if errors: return f"{context} produced error outputs but no ERROR-status job found: {errors}" if job_status == "COMPLETE": return f"{context} completed successfully" return f"{context} did not complete (status: {job_status or 'unknown'})" def _diag_print(label, upc, fields_ingest=None, fields_approval=None): if fields_ingest: print(f" DIAG {label} upc={upc}: workflow_type={fields_ingest['ingest_workflow_type'] or '(not found)'} datetime={fields_ingest['ingest_run_datetime'] or '(none)'} errors=[{fields_ingest['ingest_errors'] or 'none'}] errored_jobs=[{fields_ingest['ingest_errored_jobs'] or 'none'}]") if fields_approval: print(f" DIAG approval upc={upc}: errors=[{fields_approval['approval_errors'] or 'none'}] errored_jobs=[{fields_approval['approval_errored_jobs'] or 'none'}]") def main(upcs_file, ticket_id, sql_username): with open(upcs_file) as f: upcs = [line.strip() for line in f if line.strip()] print(f"Loaded {len(upcs)} UPCs") products = get_products_by_upcs(upcs) found_upcs = {str(p["upc"]) for p in products} skipped = [_skip(upc, "not found in DB") for upc in upcs if upc not in found_upcs] resolved = [] ows_conn = get_ows_video_conn() try: for p in products: product_id = p["product_id"] upc = str(p["upc"]) pipeline_run_id = p["latest_pipeline_run_id"] approval_job_id = p["latest_approval_job_id"] thumbnail_path = p["thumbnail_path"] thumbnail_at_ms = p["thumbnail_at_milliseconds"] # Case 1: both already set — split by whether approval workflow has run if thumbnail_path and thumbnail_at_ms is not None: if approval_job_id: asset_types = get_video_assets(product_id, ows_conn) has_s_images = any(t.startswith("video_image_S") for t in asset_types) if has_s_images: skipped.append(_skip(upc, "thumbnail_path, thumbnail_at_milliseconds, and latest_approval_job_id already set")) continue # No S_ images — check both ingest and approval for issues i_fields = {} if pipeline_run_id: i_fields = _ingest_fields(diagnose_pipeline_run(pipeline_run_id, ows_conn)) a_fields = _approval_fields(diagnose_pipeline_run(approval_job_id, ows_conn)) _diag_print("ingest", upc, i_fields, a_fields) reasons = ["no S_ images in video_asset"] if pipeline_run_id: reasons.append(_context_reason(i_fields["ingest_errors"], i_fields["ingest_errored_jobs"], i_fields["_ingest_job_status"], "ingest")) reasons.append(_context_reason(a_fields["approval_errors"], a_fields["approval_errored_jobs"], a_fields["_approval_job_status"], "approval workflow")) skipped.append(_skip(upc, " — ".join(reasons), **{k: v for k, v in {**i_fields, **a_fields}.items() if not k.startswith("_")})) continue # Approval workflow never ran — diagnose ingest to surface why i_fields = {} if pipeline_run_id: i_fields = _ingest_fields(diagnose_pipeline_run(pipeline_run_id, ows_conn)) _diag_print("ingest", upc, i_fields) reason = "approval workflow never ran — " + _context_reason(i_fields["ingest_errors"], i_fields["ingest_errored_jobs"], i_fields["_ingest_job_status"], "latest ingest") else: reason = "approval workflow never ran — no latest_pipeline_run_id, ingest may not have completed" skipped.append(_skip(upc, reason, **{k: v for k, v in i_fields.items() if not k.startswith("_")})) continue # Case 2: thumbnail_path set but thumbnail_at_milliseconds missing — derive from path if thumbnail_path and thumbnail_at_ms is None: derived_ms = parse_ms_from_path(thumbnail_path) if derived_ms is None: skipped.append(_skip(upc, f"thumbnail_path set but thumbnail_at_milliseconds not derivable (unexpected path format: {thumbnail_path})")) continue i_fields = {} if pipeline_run_id: i_fields = _ingest_fields(diagnose_pipeline_run(pipeline_run_id, ows_conn)) _diag_print("ingest", upc, i_fields) a_fields = {} if approval_job_id: a_fields = _approval_fields(diagnose_pipeline_run(approval_job_id, ows_conn)) _diag_print(None, upc, fields_approval=a_fields) has_errors = any([ i_fields.get("ingest_errors"), i_fields.get("ingest_errored_jobs"), a_fields.get("approval_errors"), a_fields.get("approval_errored_jobs"), ]) if has_errors: reasons = [f"thumbnail_path set, ms derivable as {derived_ms}ms"] if pipeline_run_id: reasons.append(_context_reason(i_fields["ingest_errors"], i_fields["ingest_errored_jobs"], i_fields["_ingest_job_status"], "ingest")) if approval_job_id: reasons.append(_context_reason(a_fields["approval_errors"], a_fields["approval_errored_jobs"], a_fields["_approval_job_status"], "approval workflow")) skipped.append(_skip(upc, " — ".join(reasons), **{k: v for k, v in {**i_fields, **a_fields}.items() if not k.startswith("_")})) print(f" WARN upc={upc} product={product_id}: ms derivable but pipeline errors detected — CSV only") else: print(f" MS upc={upc} product={product_id}: deriving {derived_ms}ms from existing path {thumbnail_path}") resolved.append({ "product_id": product_id, "upc": upc, "thumbnail_path": thumbnail_path, "thumbnail_at_ms": derived_ms, "set_path": False, }) continue # Case 3: no thumbnail_path — need to find a frame in S3 if not pipeline_run_id: skipped.append(_skip(upc, "no thumbnail_path and no latest_pipeline_run_id — ingest never completed")) continue # Diagnose ingest before checking S3 i_fields = _ingest_fields(diagnose_pipeline_run(pipeline_run_id, ows_conn)) _diag_print("ingest", upc, i_fields) found_path, found_ms = pick_first_frame(pipeline_run_id) if not found_path: # Always check approval job too if one exists a_fields = {} if approval_job_id: a_fields = _approval_fields(diagnose_pipeline_run(approval_job_id, ows_conn)) _diag_print(None, upc, fields_approval=a_fields) reasons = ["no JPG frames in S3"] reasons.append(_context_reason(i_fields["ingest_errors"], i_fields["ingest_errored_jobs"], i_fields["_ingest_job_status"], "ingest")) if approval_job_id: reasons.append(_context_reason(a_fields["approval_errors"], a_fields["approval_errored_jobs"], a_fields["_approval_job_status"], "approval workflow")) skipped.append(_skip(upc, " — ".join(reasons), **{k: v for k, v in {**i_fields, **a_fields}.items() if not k.startswith("_")})) print(f" SKIP upc={upc}: no frames in S3") continue # Frames found — diagnose both jobs; route to CSV if errors found, SQL if clean a_fields = {} if approval_job_id: a_fields = _approval_fields(diagnose_pipeline_run(approval_job_id, ows_conn)) _diag_print(None, upc, fields_approval=a_fields) has_errors = any([ i_fields["ingest_errors"], i_fields["ingest_errored_jobs"], a_fields.get("approval_errors"), a_fields.get("approval_errored_jobs"), ]) if has_errors: reasons = ["frames found in S3"] reasons.append(_context_reason(i_fields["ingest_errors"], i_fields["ingest_errored_jobs"], i_fields["_ingest_job_status"], "ingest")) if approval_job_id: reasons.append(_context_reason(a_fields["approval_errors"], a_fields["approval_errored_jobs"], a_fields["_approval_job_status"], "approval workflow")) skipped.append(_skip(upc, " — ".join(reasons), **{k: v for k, v in {**i_fields, **a_fields}.items() if not k.startswith("_")})) print(f" WARN upc={upc} product={product_id}: frames found but pipeline errors detected — CSV only") else: resolved.append({ "product_id": product_id, "upc": upc, "thumbnail_path": found_path, "thumbnail_at_ms": found_ms, "set_path": True, }) print(f" OK upc={upc} product={product_id}: {found_path} at {found_ms}ms") finally: ows_conn.close() run_ts = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") out_dir = os.path.join("outputs", run_ts) os.makedirs(out_dir, exist_ok=True) sql_path = os.path.join(out_dir, f"{ticket_id}_BACKFILL_video_thumbnails.sql") skipped_path = os.path.join(out_dir, "video_errors.csv") if resolved: write_sql(resolved, ticket_id, sql_username, sql_path) full = sum(1 for r in resolved if r["set_path"]) ms_only = sum(1 for r in resolved if not r["set_path"]) print(f"\nSQL written to {sql_path} ({full} full UPDATEs, {ms_only} milliseconds-only UPDATEs)") print(f"Drop it in database/ows_video/build/changelog/dml/ and open a PR.") else: print("\nNo products resolved — no SQL file written.") if skipped: skipped.sort(key=lambda r: r["ingest_run_datetime"] or "") write_skipped_csv(skipped, skipped_path) print(f"Skipped CSV written to {skipped_path} ({len(skipped)} UPCs)") if __name__ == "__main__": ticket = sys.argv[1] if len(sys.argv) > 1 else "MAINT" username = sys.argv[2] if len(sys.argv) > 2 else "" main("inputs/upcs.txt", ticket, username)