#!/usr/bin/env python3 """ analyze_atmos_deliveries.py Queries Snowflake for full-package DMS 187 deliveries (excluding metadata updates) that may have overwritten Dolby Atmos spatial audio, then checks each delivery's XML metadata in S3 for ImmersiveEdition. Usage: python analyze_atmos_deliveries.py [--date YYYY-MM-DD] python analyze_atmos_deliveries.py --all-time python analyze_atmos_deliveries.py --date YYYY-MM-DD --retrieve-glacierized-xml --date: Delivery date to check (default: yesterday) --all-time: Check all historical DMS 187 deliveries (enables Glacier restore) --retrieve-glacierized-xml: Submit expedited Glacier restore requests for archived XMLs and wait Setup: cp .env.shadow .env # fill in your credentials pip install -r requirements.txt """ from __future__ import annotations import argparse import csv import os import sys import time from datetime import date, datetime, timedelta from pathlib import Path import boto3 from botocore.exceptions import ClientError from cryptography.hazmat.primitives import serialization from dotenv import load_dotenv import snowflake.connector from snowflake.connector import DictCursor load_dotenv(Path(__file__).parent / ".env") S3_BUCKET = "prod-vector-audit" IMMERSIVE_MARKER = "ImmersiveEdition" OUTPUT_FIELDS = ["UPC", "LAST_1565_DELIVERY_DATE", "LAST_187_DELIVERY_DATE", "ATMOS_IN_DELIVERY", "STATUS", "FILENAME"] _QUERY_BASE = """ WITH atmos_delivered AS ( SELECT UPC, MAX(DATE_DELIVERED) AS last_1565_delivery_date FROM ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.DELIVERY_HISTORY WHERE CUSTOMER_MASTER_MASTER_ID = 1565 AND ENCODER_ID = 26 AND _FIVETRAN_DELETED = FALSE GROUP BY UPC ), dms_187_deliveries AS ( SELECT UPC, DATE_DELIVERED AS last_187_delivery_date FROM ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.DELIVERY_HISTORY WHERE CUSTOMER_MASTER_MASTER_ID = 187 AND _FIVETRAN_DELETED = FALSE {date_filter} QUALIFY ROW_NUMBER() OVER (PARTITION BY UPC ORDER BY DATE_DELIVERED DESC) = 1 ) SELECT d.UPC, a.last_1565_delivery_date, d.last_187_delivery_date, eqd.ENCODING_QUEUE_DETAIL_ID, 'metadata/' || MD5(CAST(eqd.ENCODING_QUEUE_DETAIL_ID AS VARCHAR)) || '-0.xml' AS filename FROM dms_187_deliveries d INNER JOIN atmos_delivered a ON a.UPC = d.UPC AND d.last_187_delivery_date > a.last_1565_delivery_date INNER JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.RELEASES r ON r.UPC = d.UPC AND r._FIVETRAN_DELETED = FALSE AND r.RELEASE_STATUS = 'in_content' INNER JOIN ORCHARD_APP_REPORTING_V2.REPORTSDD_DIRECT_DELIVERY.ENCODING_QUEUE_DETAIL eqd ON eqd.UPC = d.UPC AND eqd.DMS_MASTER_MASTER_ID = 187 AND eqd.STATUS = 'delivered' AND DATE(eqd.DELIVERY_ENDED) = d.last_187_delivery_date INNER JOIN ORCHARD_APP_REPORTING_V2.REPORTSDD_DIRECT_DELIVERY.ENCODING_QUEUE eq ON eqd.encoding_queue_id = eq.encoding_queue_id AND eq._fivetran_deleted = FALSE AND eq.meta_update = 'N' QUALIFY ROW_NUMBER() OVER (PARTITION BY d.UPC ORDER BY eqd.ENCODING_QUEUE_DETAIL_ID DESC) = 1 """ def build_query(delivery_date: str | None) -> tuple[str, dict]: if delivery_date: query = _QUERY_BASE.format(date_filter="AND DATE_DELIVERED = %(delivery_date)s") params = {"delivery_date": delivery_date} else: query = _QUERY_BASE.format(date_filter="") params = {} return query, params def get_snowflake_connection(): key_path = os.getenv("SNOWFLAKE_PRIVATE_KEY_PATH") if not key_path: sys.exit("Error: SNOWFLAKE_PRIVATE_KEY_PATH env var is required") key_file = Path(key_path).expanduser() if not key_file.exists(): sys.exit(f"Error: private key not found: {key_file}") user = os.getenv("SNOWFLAKE_USER") if not user: sys.exit("Error: SNOWFLAKE_USER env var is required") passphrase = os.getenv("SNOWFLAKE_KEY_PASSPHRASE") try: with open(key_file, "rb") as f: p_key = serialization.load_pem_private_key( f.read(), password=passphrase.encode() if passphrase else None, ) except (TypeError, ValueError) as e: sys.exit(f"Error: failed to load Snowflake private key ({key_file}): {e}") private_key_bytes = p_key.private_bytes( encoding=serialization.Encoding.DER, format=serialization.PrivateFormat.PKCS8, encryption_algorithm=serialization.NoEncryption(), ) return snowflake.connector.connect( account=os.getenv("SNOWFLAKE_ACCOUNT", "delphi.us-east-1"), user=user, warehouse=os.getenv("SNOWFLAKE_WAREHOUSE", "DEV_PERFORMANCE_WAREHOUSE"), database="ORCHARD_APP_REPORTING_V2", private_key=private_key_bytes, ) ARCHIVED = object() GLACIER_POLL_INTERVAL = 30 GLACIER_TIMEOUT = 300 GLACIER_RESTORE_CUTOFF = date(2026, 4, 17) # deliveries on or after this date are checked via XML, spatial asset ingester existed on this date def fetch_xml(s3_client, filename: str) -> str | None | object: try: response = s3_client.get_object(Bucket=S3_BUCKET, Key=filename) return response["Body"].read().decode("utf-8") except ClientError as e: code = e.response["Error"]["Code"] if code in ("NoSuchKey", "404"): print(f" [WARN] Not found: s3://{S3_BUCKET}/{filename}") return None if code == "InvalidObjectState": print(f" [ARCHIVED] s3://{S3_BUCKET}/{filename}") return ARCHIVED print(f" [ERROR] s3://{S3_BUCKET}/{filename}: {e}") return None def check_restore_status(s3_client, key: str) -> str: """Returns: 'available', 'restoring', 'archived', or 'not_found'""" try: head = s3_client.head_object(Bucket=S3_BUCKET, Key=key) restore = head.get("Restore") if restore is None: storage_class = head.get("StorageClass", "") if "GLACIER" in storage_class or "DEEP_ARCHIVE" in storage_class: return "archived" return "available" if 'ongoing-request="true"' in restore: return "restoring" if 'ongoing-request="false"' in restore: return "available" return "archived" except ClientError as e: code = e.response["Error"]["Code"] if code in ("NoSuchKey", "404"): return "not_found" raise def restore_archived_files(s3_client, keys: list[str]) -> set[str]: """Submit expedited restore requests for all keys, then poll until ready. Returns the set of keys that timed out.""" pending = set() for key in keys: try: head = s3_client.head_object(Bucket=S3_BUCKET, Key=key) storage_class = head.get("StorageClass", "") tier = "Standard" if "DEEP_ARCHIVE" in storage_class else "Expedited" s3_client.restore_object( Bucket=S3_BUCKET, Key=key, RestoreRequest={"Days": 1, "GlacierJobParameters": {"Tier": tier}}, ) print(f" [RESTORE REQUESTED] {key}") pending.add(key) except ClientError as e: if e.response["Error"]["Code"] == "RestoreAlreadyInProgress": print(f" [RESTORE ALREADY IN PROGRESS] {key}") pending.add(key) else: print(f" [ERROR] Could not request restore for {key}: {e}") elapsed = 0 while pending and elapsed < GLACIER_TIMEOUT: time.sleep(GLACIER_POLL_INTERVAL) elapsed += GLACIER_POLL_INTERVAL still_pending = {k for k in pending if check_restore_status(s3_client, k) != "available"} ready = pending - still_pending for key in ready: print(f" [READY] {key}") pending = still_pending if pending: print(f" [WAITING] {len(pending)} file(s) still restoring... ({elapsed}s elapsed)") if pending: print(f" [TIMEOUT] {len(pending)} file(s) did not restore within {GLACIER_TIMEOUT}s") return pending def _parse_date(value: str) -> str: try: datetime.strptime(value, "%Y-%m-%d") except ValueError: raise argparse.ArgumentTypeError(f"Invalid date '{value}'; expected YYYY-MM-DD") return value def main(): parser = argparse.ArgumentParser( description="Check DMS 187 deliveries for Dolby Atmos presence in delivery XML" ) group = parser.add_mutually_exclusive_group() group.add_argument("--date", type=_parse_date, help="Delivery date to check (YYYY-MM-DD, default: yesterday)") group.add_argument("--all-time", action="store_true", help="Check all historical DMS 187 deliveries") parser.add_argument("--retrieve-glacierized-xml", action="store_true", help="Submit expedited Glacier restore requests for archived XMLs and wait") args = parser.parse_args() retrieve_glacierized = args.retrieve_glacierized_xml or args.all_time if args.all_time: delivery_date = None print("Checking all historical DMS 187 deliveries") else: delivery_date = args.date or str(date.today() - timedelta(days=1)) print(f"Checking DMS 187 deliveries for date: {delivery_date}") print("Connecting to Snowflake...") conn = get_snowflake_connection() cursor = conn.cursor(DictCursor) query, params = build_query(delivery_date) print("Fetching deliveries...") cursor.execute(query, params) rows = cursor.fetchall() cursor.close() conn.close() print(f"Found {len(rows)} deliveries to check.\n") if not rows: print("Nothing to check.") return s3_client = boto3.client("s3") results = [] timed_out_keys: set[str] = set() if retrieve_glacierized: glacier_keys = [ row["FILENAME"] for row in rows if row["LAST_187_DELIVERY_DATE"] >= GLACIER_RESTORE_CUTOFF and check_restore_status(s3_client, row["FILENAME"]) in ("archived", "restoring") ] if glacier_keys: print(f"\nFound {len(glacier_keys)} archived/restoring file(s) after cutoff — submitting expedited restores...") timed_out_keys = restore_archived_files(s3_client, glacier_keys) print() for row in rows: upc = row["UPC"] filename = row["FILENAME"] delivery_date_val = row["LAST_187_DELIVERY_DATE"] print(f"UPC={upc} file={filename}") if delivery_date_val < GLACIER_RESTORE_CUTOFF: atmos_in_delivery = False status = "No Atmos — needs correction (pre-cutoff)" else: xml = fetch_xml(s3_client, filename) if xml is ARCHIVED: atmos_in_delivery = None status = "Archived — restore timed out" if filename in timed_out_keys else "Archived — cannot check" elif xml is None: atmos_in_delivery = None status = "XML not found" else: atmos_in_delivery = IMMERSIVE_MARKER in xml status = "Atmos present" if atmos_in_delivery else "No Atmos — needs correction" print(f" [{status}]") results.append({ "UPC": row["UPC"], "LAST_1565_DELIVERY_DATE": row["LAST_1565_DELIVERY_DATE"], "LAST_187_DELIVERY_DATE": row["LAST_187_DELIVERY_DATE"], "ATMOS_IN_DELIVERY": atmos_in_delivery, "STATUS": status, "FILENAME": filename, }) timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") label = "all_time" if delivery_date is None else delivery_date output_filename = f"atmos_check_{label}_{timestamp}.csv" output_path = Path(__file__).parent / "outputs" / output_filename output_path.parent.mkdir(parents=True, exist_ok=True) with open(output_path, "w", newline="", encoding="utf-8") as f: writer = csv.DictWriter(f, fieldnames=OUTPUT_FIELDS) writer.writeheader() writer.writerows(results) needs_correction = [r for r in results if "needs correction" in r["STATUS"]] atmos_confirmed = [r for r in results if r["STATUS"] == "Atmos present"] archived = [r for r in results if r["STATUS"] == "Archived — cannot check"] xml_missing = [r for r in results if r["STATUS"] == "XML not found"] print(f"\nResults written to: {output_path}") print(f"Needs correction: {len(needs_correction)}") print(f"Atmos confirmed in XML: {len(atmos_confirmed)}") print(f"Archived — cannot check: {len(archived)}") print(f"XML not found: {len(xml_missing)}") if __name__ == "__main__": main()