from concurrent.futures import ProcessPoolExecutor from concurrent.futures.process import BrokenProcessPool from multiprocessing import get_context from typing import Final, cast from pydantic import AnyUrl from src.atmos.alignment_check import apply_alignment_measurement from src.atmos.lfe_check import apply_lfe_measurement from src.atmos.loudness_check import apply_loudness_measurement from src.atmos.mediainfo_check import check_mediainfo, run_mediainfo from src.atmos.models import ( AtmosValidationBuilder, AtmosValidationFindingCode, AtmosValidationRequest, AtmosValidationResult, ) from src.atmos.render_checks import run_render_checks from src.atmos.render_errors import RenderWorkerError from src.atmos.silent_object_check import apply_silent_object_measurement from src.clients import ows_assets, s3 from src.config import presigned_url_ttl_seconds from src.worker import run_handler # EAR rendering is CPU-bound and memory-heavy, so the render-based checks run in a separate process: # it isolates the render's ~2.2 GB peak from the long-running poller and keeps a render crash from # terminating the worker. spawn, not fork: boto3 clients aren't fork-safe; spawn gives each worker a # fresh interpreter. _MAX_RENDER_WORKERS: Final = 1 # one render worker; both layouts render from a single shared ADM parse in one pass _RENDER_POOL_CONTEXT: Final = get_context("spawn") def validate(atmos_url: AnyUrl, stereo_url: AnyUrl) -> AtmosValidationBuilder: atmos_validation_builder = AtmosValidationBuilder() check_mediainfo(atmos_url, atmos_validation_builder) check_stereo_reference(stereo_url, atmos_validation_builder) return atmos_validation_builder def check_stereo_reference( stereo_url: AnyUrl, atmos_validation_builder: AtmosValidationBuilder, ) -> None: stereo_output = run_mediainfo(stereo_url) atmos_validation_builder.update_metadata( stereo_reference_is_truncated=stereo_output.is_truncated, stereo_reference_duration_ms=stereo_output.duration_ms, ) if stereo_output.is_truncated: atmos_validation_builder.error( AtmosValidationFindingCode.STEREO_FILE_TRUNCATED, "Stereo reference file is truncated — actual size is less than the size declared in the header", ) return _check_duration_match(atmos_validation_builder) def _check_duration_match(atmos_validation_builder: AtmosValidationBuilder) -> None: metadata = atmos_validation_builder.metadata if metadata is None or metadata.is_truncated: return duration_diff_ms = metadata.stereo_reference_duration_diff_ms if duration_diff_ms is None: return threshold_ms: Final = 2000 if duration_diff_ms > threshold_ms: atmos_validation_builder.error( AtmosValidationFindingCode.DURATION_MISMATCH, f"Atmos duration ({metadata.duration_ms} ms) and stereo reference " f"({metadata.stereo_reference_duration_ms} ms) differ by {duration_diff_ms} ms, " f"must not exceed {threshold_ms} ms", ) def _run_render_checks( atmos_bucket: str, atmos_key: str, stereo_reference_bucket: str, stereo_reference_key: str, atmos_validation_builder: AtmosValidationBuilder, ) -> None: # Advisory render-based checks run only on an otherwise-valid delivery: a structural error # already fails it, and rendering a malformed master is wasted work. The single render worker # produces every render-based measurement (loudness, alignment, LFE, silent-object) from one # shared ADM parse. if atmos_validation_builder.errors: return with ProcessPoolExecutor(max_workers=_MAX_RENDER_WORKERS, mp_context=_RENDER_POOL_CONTEXT) as executor: render_future = executor.submit( run_render_checks, atmos_bucket, atmos_key, stereo_reference_bucket, stereo_reference_key, ) try: render_check_measurements = render_future.result() except BrokenProcessPool as e: raise RenderWorkerError( "render worker terminated abruptly (likely out of memory during the render); " f"check the task memory limit against max_workers={_MAX_RENDER_WORKERS}" ) from e apply_loudness_measurement(render_check_measurements.loudness, atmos_validation_builder) apply_alignment_measurement(render_check_measurements.alignment, atmos_validation_builder) apply_lfe_measurement(render_check_measurements.lfe, atmos_validation_builder) apply_silent_object_measurement(render_check_measurements.silent_object, atmos_validation_builder) def _validate(atmos_validation_request: AtmosValidationRequest) -> AtmosValidationResult: stereo_reference_bucket, stereo_reference_key = _resolve_stereo_reference(atmos_validation_request) url_expires_in_seconds = presigned_url_ttl_seconds() atmos_url = s3.create_presigned_url( atmos_validation_request.atmos_bucket, atmos_validation_request.atmos_key, expires_in_seconds=url_expires_in_seconds, ) stereo_url = s3.create_presigned_url( stereo_reference_bucket, stereo_reference_key, expires_in_seconds=url_expires_in_seconds, ) atmos_validation_builder = validate(atmos_url, stereo_url) _run_render_checks( atmos_validation_request.atmos_bucket, atmos_validation_request.atmos_key, stereo_reference_bucket, stereo_reference_key, atmos_validation_builder, ) atmos_validation_result = atmos_validation_builder.to_result() if atmos_validation_request.report_validation_result: ows_assets.post_validation_result(atmos_validation_request.atmos_key, atmos_validation_result) return atmos_validation_result def _resolve_stereo_reference(atmos_validation_request: AtmosValidationRequest) -> tuple[str, str]: if atmos_validation_request.lookup_stereo_reference: return ows_assets.get_stereo_for_asset(atmos_validation_request.atmos_key) return ( cast(str, atmos_validation_request.stereo_reference_bucket), cast(str, atmos_validation_request.stereo_reference_key), ) def main() -> None: run_handler(AtmosValidationRequest, _validate) if __name__ == "__main__": main()