"""Spatial audio asset ingester. For a given stereo/spatial UPC+ISRC pair, this script: 1. Resolves the stereo release_id (via OWS Product) and track_id (via OWS Track). 2. Locates the spatial .wav file in the Atmos source S3 bucket. 3. Uploads the spatial asset to ows-assets (reads from S3, creates asset_upload record, and copies the file) via python-ows-assets-uploader. 4. Creates a track_spatial record in OWS Track linking stereo track -> spatial ISRC. Accepts either a CSV file path via the BULK_SPATIAL_FILE environment variable, or individual inputs via environment variables (STEREO_UPC, STEREO_ISRC, SPATIAL_UPC, SPATIAL_ISRC). ENVIRONMENT is always read from the environment variable of the same name. In batch mode, all rows are attempted regardless of per-row failures. Exits with code 1 if any row failed. """ import csv import os import sys import time from concurrent.futures import ThreadPoolExecutor, as_completed from http import HTTPStatus from typing import Any import httpx import ows_assets_uploader from owsclient import M2MTokenManager, OwsClient from src import config from src.connectors import art_relations, ows_product, ows_product_digital, ows_track, s3 from src.exceptions import DisplayUpcNotFound client = OwsClient( environment=config.ENVIRONMENT, service_name=config.SERVICE_NAME, timeout=config.DEFAULT_OWS_CLIENT_TIMEOUT, m2m_token_manager=M2MTokenManager( secrets_manager=config.secrets_manager, environment=config.TOKEN_MANAGER_ENVIRONMENT, service_name=config.SERVICE_NAME, ), ) s3_connector = s3.S3Connector(config.ENVIRONMENT) ows_product_connector = ows_product.OwsProductConnector(client) ows_product_digital_connector = ows_product_digital.OwsProductDigitalConnector(client) ows_track_connector = ows_track.OwsTrackConnector(client) art_relations_connector = art_relations.ArtRelationsConnector( host=config.ART_RELATIONS_DB_HOST, user=config.ART_RELATIONS_DB_USER, password=config.ART_RELATIONS_DB_PASSWORD, database=config.ART_RELATIONS_DB_NAME, ) REQUIRED_COLUMNS = {"STEREO_UPC", "STEREO_ISRC", "SPATIAL_UPC", "SPATIAL_ISRC"} ATMOS_ASSET_UPLOAD_TYPE = "atmos" DEFAULT_BULK_MAX_WORKERS = 2 OWS_ASSETS_404_NUM_RETRIES = 5 OWS_ASSETS_404_INITIAL_BACKOFF_SECONDS = 0.5 def _upload_with_404_retry(**kwargs: Any) -> str: """Call ows_assets_uploader.upload, retrying on transient 404s with exponential backoff. ows-assets occasionally returns 404 on calls that read the asset_upload record immediately after it was created, due to read-after-write consistency lag. The retry restarts the upload from scratch (new POST), which is cheap because the 404 happens early — before any file bytes are uploaded. """ delay = OWS_ASSETS_404_INITIAL_BACKOFF_SECONDS for attempt in range(1, OWS_ASSETS_404_NUM_RETRIES + 1): try: return str(ows_assets_uploader.upload(**kwargs)) except httpx.HTTPStatusError as e: if e.response.status_code != HTTPStatus.NOT_FOUND or attempt == OWS_ASSETS_404_NUM_RETRIES: raise print(f"ows-assets returned 404; retrying upload in {delay:.1f}s") time.sleep(delay) delay *= 2 raise RuntimeError("unreachable") def _run_row(row: dict[str, str], upc_map: dict[str, int]) -> None: stereo_display_upc = row["STEREO_UPC"] stereo_isrc = row["STEREO_ISRC"] spatial_display_upc = row["SPATIAL_UPC"] spatial_isrc = row["SPATIAL_ISRC"] if stereo_display_upc not in upc_map: raise DisplayUpcNotFound(f"No release found for stereo UPC {stereo_display_upc!r}") if spatial_display_upc not in upc_map: raise DisplayUpcNotFound(f"No release found for spatial UPC {spatial_display_upc!r}") print(f"\n[{spatial_display_upc}/{spatial_isrc}] Starting ingestion...") ingest(upc_map[stereo_display_upc], stereo_isrc, upc_map[spatial_display_upc], spatial_display_upc, spatial_isrc) print(f"[{spatial_display_upc}/{spatial_isrc}] Ingestion complete") def main() -> None: """Entrypoint.""" bulk_file = os.environ.get("BULK_SPATIAL_FILE") if bulk_file: with open(bulk_file, newline="") as f: reader = csv.DictReader(f) if not REQUIRED_COLUMNS.issubset(set(reader.fieldnames or [])): missing = REQUIRED_COLUMNS - set(reader.fieldnames or []) raise ValueError(f"input.csv is missing required columns: {', '.join(sorted(missing))}") rows = list(reader) else: row = { "STEREO_UPC": os.environ["STEREO_UPC"], "STEREO_ISRC": os.environ["STEREO_ISRC"], "SPATIAL_UPC": os.environ["SPATIAL_UPC"], "SPATIAL_ISRC": os.environ["SPATIAL_ISRC"], } try: upc_map = art_relations_connector.get_upcs_by_display_upcs({row["STEREO_UPC"], row["SPATIAL_UPC"]}) _run_row(row, upc_map) except Exception as e: print(f"[{row['SPATIAL_UPC']}/{row['SPATIAL_ISRC']}] Error: {e}") sys.exit(getattr(e, "exit_code", 1)) return all_display_upcs = {upc for row in rows for upc in [row["STEREO_UPC"], row["SPATIAL_UPC"]]} try: upc_map = art_relations_connector.get_upcs_by_display_upcs(all_display_upcs) except Exception as e: print(f"Error fetching UPCs from art_relations: {e}") sys.exit(1) max_workers = int(os.environ.get("BULK_MAX_WORKERS", DEFAULT_BULK_MAX_WORKERS)) print(f"Processing {len(rows)} row(s) with {max_workers} workers...") failures = [] with ThreadPoolExecutor(max_workers=max_workers) as executor: future_to_row = {executor.submit(_run_row, row, upc_map): row for row in rows} for future in as_completed(future_to_row): row = future_to_row[future] try: future.result() except Exception as e: spatial_display_upc = row["SPATIAL_UPC"] spatial_isrc = row["SPATIAL_ISRC"] print(f"[{spatial_display_upc}/{spatial_isrc}] Error: {e}") exit_code = getattr(e, "exit_code", 1) failures.append((spatial_display_upc, spatial_isrc, str(e), exit_code)) print(f"\n{len(rows) - len(failures)}/{len(rows)} row(s) completed successfully.") if failures: print(f"{len(failures)} row(s) failed:") for spatial_display_upc, spatial_isrc, error, exit_code in failures: print(f" [{spatial_display_upc}/{spatial_isrc}] (exit code {exit_code}) {error}") sys.exit(1) def ingest( stereo_upc: int, stereo_isrc: str, spatial_upc: int, spatial_display_upc: str, spatial_isrc: str, ) -> None: """Run the full ingestion pipeline for a single stereo/spatial pair.""" ######### Fetch stereo release_id and track_id ######### print(f"\n[{spatial_display_upc}/{spatial_isrc}] Fetching stereo track information...") stereo_release_id = ows_product_connector.get_release_id_by_upc(stereo_upc) stereo_track_id = ows_track_connector.get_track_id_by_upc_and_isrc(stereo_upc, stereo_isrc) print( f"[{spatial_display_upc}/{spatial_isrc}] Fetched stereo release_id={stereo_release_id} and track_id={stereo_track_id}" ) ######### Fetch the spatial asset in S3 ######### print(f"\n[{spatial_display_upc}/{spatial_isrc}] Searching for spatial asset in input bucket...") source_key = s3_connector.find_spatial_asset_key(spatial_display_upc, spatial_isrc) source_s3_url = f"s3://{s3_connector.source_bucket}/{source_key}" print(f"[{spatial_display_upc}/{spatial_isrc}] Found spatial asset at {source_s3_url}") ######### Create spatial asset upload ######### print(f"\n[{spatial_display_upc}/{spatial_isrc}] Uploading spatial asset...") destination_filename = _upload_with_404_retry( source_file_location=source_s3_url, product_id=stereo_release_id, track_id=stereo_track_id, asset_upload_type=ATMOS_ASSET_UPLOAD_TYPE, ows_client=client, ) print(f"[{spatial_display_upc}/{spatial_isrc}] Uploaded spatial asset, destination_filename={destination_filename}") ######### Create track_spatial record ######### print(f"\n[{spatial_display_upc}/{spatial_isrc}] Creating spatial track record...") ows_track_connector.create_track_spatial( track_id=stereo_track_id, spatial_isrc=spatial_isrc, ) print(f"[{spatial_display_upc}/{spatial_isrc}] Created track_spatial record") ######### Create release_spatial record ######### print(f"\n[{spatial_display_upc}/{spatial_isrc}] Creating spatial release record...") ows_product_digital_connector.create_product_spatial( stereo_release_id=stereo_release_id, spatial_upc=spatial_upc, ) print(f"[{spatial_display_upc}/{spatial_isrc}] Created release_spatial record") if __name__ == "__main__": main()