"""Job logic tests.""" import json import uuid from collections.abc import Generator from typing import Any from unittest.mock import MagicMock, call import jsonschema import pytest from pytest_mock import MockerFixture from spec.jsonschema import jobs as jobs_validation from video import api from video.connectors import mediaconvert, mysql, stepfunctions from video.constants import job_io_fields, job_statuses, job_types from video.exceptions import InvalidRequest from video.logic import job as job_logic from video.models.s3 import cloud_transfer_progress from video.models.sql import queries from video.models.sql.classes import product_video def mock_persist(mocker: MockerFixture, get_job_response: Any) -> MagicMock: """Mock the persist function.""" persist_mock = mocker.patch.object( queries, "persist_job_data", autospec=True, ) mocker.patch.object( queries, "get_jobs", return_value=[get_job_response], autospec=True, ) return persist_mock def test_manage_job(mocker: MockerFixture) -> None: """Test manage job.""" get_job_response = {"applesauce": "bananas", "id": 3} persist_mock = mock_persist(mocker, get_job_response) job_data = {"applesauce": "bananas"} manage_job_response = job_logic.persist_new_job_data(job_data) persist_mock.assert_called_once_with(job_data) assert get_job_response == manage_job_response session_mock = MagicMock() @pytest.mark.parametrize( ( "test_description", "data", "manage_job_response", "expected_mediaconvert_cancel_job_calls", "expected_stepfunctions_stop_execution_calls", "expected_get_jobs_calls", "get_jobs_responses", "expected_persist_job_data_calls", ), [ ( "Cancel non-workflow job with no parent.", {"id": 123, "status": job_statuses.CANCELLED}, { "id": 123, "status": job_statuses.CANCELLED, "type": "some_non_workflow_job_type", }, [], [], [], [], [], ), ( "Error non-workflow job with no parent.", {"id": 123, "status": job_statuses.ERROR}, { "id": 123, "status": job_statuses.ERROR, "type": "some_non_workflow_job_type", }, [], [], [], [], [], ), ( "Cancel non-workflow job with a parent.", {"id": 123, "status": job_statuses.CANCELLED}, { "id": 123, "parent_id": 7, "status": job_statuses.CANCELLED, "type": "some_non_workflow_job_type", }, [], [], [], [], [], ), ( "Error non-workflow job with a parent.", {"id": 123, "status": job_statuses.ERROR}, { "id": 123, "parent_id": 7, "status": job_statuses.ERROR, "type": "some_non_workflow_job_type", }, [], [], [call({"id": 7}, session_mock)], [ [ { "id": 7, "status": job_statuses.PROGRESSING, "type": "some_non_workflow_job_type", } ] ], [], ), ( "Cancel workflow job with no parent.", {"id": 123, "status": job_statuses.CANCELLED}, { "id": 123, "status": job_statuses.CANCELLED, "type": job_types.WORKFLOW_INGEST_FROM_BROWSER, "inputs": { job_io_fields.WORKFLOW_STATE_MACHINE_EXECUTION_ARN: "smaaaaaaaarrrrrrn" }, }, [call(Id="mediconid")], [call(executionArn="smaaaaaaaarrrrrrn")], [call({"parent_id": 123}, session_mock)], [ [ {"id": 333, "status": job_statuses.COMPLETE, "inputs": None}, { "id": 334, "status": job_statuses.PROGRESSING, "inputs": { job_io_fields.CREATE_STREAMABLE_PREVIEW_AND_THUMBNAILS_MEDIACONVERT_JOB_ID: "mediconid" }, }, {"id": 335, "status": job_statuses.SUBMITTED, "inputs": None}, ] ], [ call( [ {"id": 334, "status": job_statuses.CANCELLED}, {"id": 335, "status": job_statuses.CANCELLED}, ], session=session_mock, ) ], ), ( "Error workflow job with no parent.", {"id": 123, "status": job_statuses.ERROR}, { "id": 123, "status": job_statuses.ERROR, "type": job_types.WORKFLOW_INGEST_FROM_BROWSER, }, [], [], [], [], [], ), ( "Cancel workflow job with a parent.", {"id": 123, "status": job_statuses.CANCELLED}, { "id": 123, "parent_id": 7, "status": job_statuses.CANCELLED, "type": job_types.WORKFLOW_INGEST_FROM_BROWSER, }, [], [], [], [], [], ), ( "Error workflow job with a parent.", {"id": 123, "status": job_statuses.ERROR}, {"id": 123, "parent_id": 7, "status": job_statuses.ERROR}, [call(Id="mediconid")], [], [call({"parent_id": 7}, session_mock)], [ [ { "id": 7, "status": job_statuses.PROGRESSING, "type": job_types.WORKFLOW_INGEST_FROM_BROWSER, } ], [ {"id": 333, "status": job_statuses.COMPLETE, "inputs": None}, { "id": 334, "status": job_statuses.PROGRESSING, "inputs": { job_io_fields.CREATE_STREAMABLE_PREVIEW_AND_THUMBNAILS_MEDIACONVERT_JOB_ID: "mediconid" }, }, {"id": 335, "status": job_statuses.SUBMITTED, "inputs": None}, ], ], [ call( [ {"id": 7, "status": job_statuses.ERROR}, {"id": 334, "status": job_statuses.CANCELLED}, {"id": 335, "status": job_statuses.CANCELLED}, ], session=session_mock, ) ], ), ], ) def test_if_in_workflow_apply_updates( mocker: MockerFixture, test_description: Any, data: Any, manage_job_response: Any, expected_mediaconvert_cancel_job_calls: Any, expected_stepfunctions_stop_execution_calls: Any, expected_get_jobs_calls: Any, get_jobs_responses: Any, expected_persist_job_data_calls: Any, ) -> None: """Test if_in_workflow_apply_updates.""" from contextlib import contextmanager @contextmanager def fake_db_session(**kwargs: Any) -> Generator[MagicMock, None, None]: yield session_mock mocker.patch.object(mysql, "db_session", fake_db_session) mock_mc: MagicMock = mocker.patch.object( mediaconvert, "get_mediaconvert_client", autospec=True ) mock_sf: MagicMock = mocker.patch.object( stepfunctions, "get_stepfunctions_client", autospec=True ) get_jobs_mock = mocker.patch.object( queries, "get_jobs", autospec=True, side_effect=get_jobs_responses ) persist_mock = mocker.patch.object(queries, "persist_job_data", autospec=True) job_logic.if_in_workflow_apply_updates(data, manage_job_response) mock_mc.return_value.cancel_job.assert_has_calls( expected_mediaconvert_cancel_job_calls ) mock_sf.return_value.stop_execution.assert_has_calls( expected_stepfunctions_stop_execution_calls ) get_jobs_mock.assert_has_calls(expected_get_jobs_calls) persist_mock.assert_has_calls(expected_persist_job_data_calls) @pytest.mark.parametrize( ("test_description", "existing_job", "data", "expected_schema"), [ ( "test create job validation", None, {"type": job_types.TRANSFER_FROM_BROWSER_TO_S3}, jobs_validation.build_request_schema_create_job( job_types.TRANSFER_FROM_BROWSER_TO_S3 ), ), ( "test change status validation", { "id": 3, "status": job_statuses.SETUP, "type": job_types.TRANSFER_FROM_BROWSER_TO_S3, }, {"id": 3, "status": job_statuses.SUBMITTED}, jobs_validation.build_request_schema_change_job_status( job_type=job_types.TRANSFER_FROM_BROWSER_TO_S3, current_job_status=job_statuses.SETUP, ), ), ( "test create inputs validation", { "id": 3, "status": job_statuses.SETUP, "type": job_types.TRANSFER_FROM_BROWSER_TO_S3, }, {"id": 3, "inputs": {"applesauce": "bananas"}}, jobs_validation.build_request_schema_create_job_inputs( job_types.TRANSFER_FROM_BROWSER_TO_S3 ), ), ( "test create outputs validation", { "id": 3, "status": job_statuses.PROGRESSING, "type": job_types.TRANSFER_FROM_BROWSER_TO_S3, }, {"id": 3, "outputs": {"applesauce": "bananas"}}, jobs_validation.build_request_schema_create_job_outputs( job_types.TRANSFER_FROM_BROWSER_TO_S3 ), ), ], ) def test_validate_job_management_request_data( mocker: MockerFixture, test_description: Any, existing_job: Any, data: Any, expected_schema: Any, ) -> None: """Test validate_job_management_request_data.""" validate_mock = mocker.patch.object( jsonschema, "validate", autospec=True, ) job_logic.validate_job_management_request_data(data, existing_job) validate_mock.assert_called_once_with(data, expected_schema) def test_validate_job_management_request_data_missing_type_raises() -> None: """Test that creating a job without a type raises InvalidRequest.""" from video.exceptions import InvalidRequest with pytest.raises(InvalidRequest): job_logic.validate_job_management_request_data({}, None) def test_validate_job_management_request_data_unknown_type_raises() -> None: """Test that creating a job with an unrecognized type raises InvalidRequest.""" from video.exceptions import InvalidRequest with pytest.raises(InvalidRequest): job_logic.validate_job_management_request_data( {"type": "not_a_real_type"}, None ) def test_get_jobs(mocker: MockerFixture) -> None: """Test get_jobs.""" get_jobs_mock = mocker.patch.object( queries, "get_jobs", return_value=[{"f1": "f1"}] ) get_jobs_response = job_logic.get_jobs({"f1": "f1", "f2": "f2"}) get_jobs_mock.assert_called_once_with({"f1": "f1", "f2": "f2"}) assert get_jobs_response == [{"f1": "f1"}] def test_get_file_attributes_success() -> None: """Test get_file_attributes success case.""" response = job_logic.get_file_attributes({"name": "foo.mov", "size": 2000}) assert response == {"filename": "foo.mov", "filesize": 2000} def test_get_file_attributes_raises_validation_error() -> None: """Test get_file_attributes validation failure.""" with pytest.raises(jsonschema.ValidationError) as execinfo: job_logic.get_file_attributes({"name": "bar.mov", "size": "2001"}) assert execinfo.value.validator is not None assert execinfo.value.validator_value is not None assert execinfo.value.message is not None jobs_to_setup = [ {"parent_id": 12, "type": "extract_metadata"}, {"parent_id": 12, "type": "validate_metadata"}, {"parent_id": 12, "type": "calculate_audio_stats"}, {"parent_id": 12, "type": "detect_black_intervals_at_beginning_and_end_of_video"}, {"parent_id": 12, "type": "detect_max_volume"}, {"parent_id": 12, "type": "detect_silent_intervals_at_beginning_and_end_of_audio"}, { "parent_id": 12, "type": "detect_video_location_and_dimensions_within_black_border", }, {"parent_id": 12, "type": "validate_analysis"}, {"parent_id": 12, "type": "create_streamable_preview_and_thumbnails"}, ] @pytest.mark.parametrize( ( "expected_jobs_to_setup", "setup_ingest_workflow_function", "transfer_job_type", "workflow_job_type", "workflow_job_inputs", "transfer_jobs_to_set_up", "additional_stepfunctions_inputs", "workflow_handler_inputs", ), [ ( jobs_to_setup, job_logic.setup_ingest_from_browser_workflow, job_types.TRANSFER_FROM_BROWSER_TO_S3, job_types.WORKFLOW_INGEST_FROM_BROWSER, {}, [ { "parent_id": 12, "type": job_types.TRANSFER_FROM_BROWSER_TO_S3, "inputs": { "fake_token_key": "fake_token_val", }, }, ], {}, {}, ), ( jobs_to_setup, job_logic.setup_ingest_from_browser_workflow, job_types.TRANSFER_FROM_BROWSER_TO_S3, job_types.WORKFLOW_INGEST_FROM_BROWSER, {}, [ { "parent_id": 12, "type": job_types.TRANSFER_FROM_BROWSER_TO_S3, "inputs": { "fake_token_key": "fake_token_val", }, }, ], {}, {}, ), ( jobs_to_setup, job_logic.setup_ingest_from_s3_workflow, job_types.TRANSFER_FROM_S3_TO_S3, job_types.WORKFLOW_INGEST_FROM_S3, {}, [ { "parent_id": 12, "type": job_types.TRANSFER_FROM_S3_TO_S3, "inputs": { "input_video_s3_bucket": "test-orcd-video-assets", "input_video_s3_key": "ingest/incoming/video_to_ingest.mp4", "workflow_job_id": 12, }, }, ], {}, { "path": "ingest/incoming/video_to_ingest.mp4", }, ), ( jobs_to_setup, job_logic.setup_ingest_from_s3_workflow, job_types.TRANSFER_FROM_S3_TO_S3, job_types.WORKFLOW_INGEST_FROM_S3, {}, [ { "parent_id": 12, "type": job_types.TRANSFER_FROM_S3_TO_S3, "inputs": { "input_video_s3_bucket": "test-orcd-video-assets", "input_video_s3_key": "ingest/incoming/video_to_ingest.mp4", "workflow_job_id": 12, }, }, ], {}, { "path": "ingest/incoming/video_to_ingest.mp4", }, ), ( jobs_to_setup, job_logic.setup_ingest_from_google_drive_workflow, job_types.TRANSFER_FROM_GOOGLE_DRIVE_TO_S3, job_types.WORKFLOW_INGEST_FROM_GOOGLE_DRIVE, { job_io_fields.GOOGLE_DRIVE_FILES: {"some": "files"}, job_io_fields.GOOGLE_DRIVE_AUTHORIZATION: {"some": "auth"}, }, [ { "parent_id": 12, "type": job_types.TRANSFER_FROM_GOOGLE_DRIVE_TO_S3, }, ], { job_io_fields.GOOGLE_DRIVE_FILES: {"some": "files"}, job_io_fields.GOOGLE_DRIVE_AUTHORIZATION: {"some": "auth"}, }, { "inputs": { job_io_fields.GOOGLE_DRIVE_FILES: {"some": "files"}, job_io_fields.GOOGLE_DRIVE_AUTHORIZATION: {"some": "auth"}, }, }, ), ( jobs_to_setup, job_logic.setup_ingest_from_dropbox_workflow, job_types.TRANSFER_FROM_DROPBOX_TO_S3, job_types.WORKFLOW_INGEST_FROM_DROPBOX, { job_io_fields.DROPBOX_FILES: {"some": "files"}, }, [ { "parent_id": 12, "type": job_types.TRANSFER_FROM_DROPBOX_TO_S3, }, ], { job_io_fields.DROPBOX_FILES: {"some": "files"}, }, { "inputs": { job_io_fields.DROPBOX_FILES: {"some": "files"}, }, }, ), ], ) def test_setup_ingest_workflow( mocker: MockerFixture, expected_jobs_to_setup: Any, setup_ingest_workflow_function: Any, transfer_job_type: Any, workflow_job_type: Any, workflow_job_inputs: Any, transfer_jobs_to_set_up: Any, additional_stepfunctions_inputs: Any, workflow_handler_inputs: Any, ) -> None: """Test setup ingest workflow.""" with api.app.test_request_context(): parent_id = 12 persist_mock = mocker.patch.object( queries, "persist_job_data", autospec=True, return_value=[parent_id] ) mocker.patch.object( jsonschema, "validate", autospec=True, ) mocker.patch.object( job_logic, "get_transfer_from_browser_to_s3_inputs", autospec=True, return_value={"fake_token_key": "fake_token_val"}, ) get_jobs_mock = mocker.patch.object( queries, "get_jobs", autospec=True, return_value=[ { "id": parent_id, "type": transfer_job_type, } ], ) mocker.patch.object(job_logic, "check_has_access", autospec=True) start_execution_mock = mocker.patch.object( stepfunctions, "start_execution", autospec=True, return_value={"executionArn": "cool-arn"}, ) mocker.patch.object(mysql, "db_session", autospec=True) mocker.patch.object(uuid, "uuid1", return_value="test_uuid", autospec=True) setup_ingest_workflow_function(workflow_handler_inputs) get_jobs_mock.assert_called_once_with({"parent_id": parent_id}) persist_mock.assert_has_calls( [ call( { "type": workflow_job_type, "inputs": workflow_job_inputs, }, session=mysql.db_session().__enter__(), ), call( [*transfer_jobs_to_set_up, *expected_jobs_to_setup], session=mysql.db_session().__enter__(), ), call( { "id": 12, "inputs": {"workflow_state_machine_execution_arn": "cool-arn"}, }, session=mysql.db_session().__enter__(), ), ] ) stepfunctions_input = { **additional_stepfunctions_inputs, "workflow_job_type": workflow_job_type, "workflow_job_id": parent_id, } start_execution_mock.assert_called_once_with( json.dumps(stepfunctions_input), "{}_{}".format("test_uuid", parent_id), workflow_job_type, ) @pytest.mark.parametrize( ("test_description", "additional_inputs"), [ ( "enabled is_new_video_approval_workflow_enabled", [ { "parent_id": 13, "type": "write_mezz_video_locations_to_video_asset_table", }, {"parent_id": 13, "type": "convert_thumbnails_to_tiffs"}, { "parent_id": 13, "type": "write_thumbnail_locations_to_video_asset_table", }, {"parent_id": 13, "type": "mark_video_product_as_approved"}, ], ), ], ) def test_setup_approval_workflow( mocker: MockerFixture, test_description: Any, additional_inputs: Any ) -> None: """Test setup_approval_workflow.""" parent_id = 13 persist_mock = mocker.patch.object( queries, "persist_job_data", autospec=True, return_value=[parent_id] ) mocker.patch.object(product_video, "upsert", autospec=True, return_value=parent_id) mocker.patch.object( jsonschema, "validate", autospec=True, ) get_jobs_mock = mocker.patch.object( queries, "get_jobs", autospec=True, return_value=[ { "id": parent_id, "type": job_types.WORKFLOW_APPROVAL, } ], ) mocker.patch.object(job_logic, "check_has_access", autospec=True) start_execution_mock = mocker.patch.object( stepfunctions, "start_execution", autospec=True, return_value={"executionArn": "cool-arn"}, ) mocker.patch.object(mysql, "db_session", autospec=True) mocker.patch.object(uuid, "uuid1", return_value="test_uuid", autospec=True) job_logic.setup_approval_workflow( { "workflow_ingest_job_id": 88997690, "context": {"product_id": 123, "upc": 1234}, } ) jobs_to_persist = [ {"parent_id": 13, "type": job_types.GET_CREATE_MEZZANINES_INPUTS}, {"parent_id": 13, "type": job_types.GET_PRODUCT_METADATA}, {"parent_id": 13, "type": job_types.EXTRACT_PRORES_MEZZANINE_METADATA}, {"parent_id": 13, "type": job_types.CREATE_MEZZANINES}, *additional_inputs, ] get_jobs_mock.assert_called_once_with({"parent_id": parent_id}) persist_mock.assert_has_calls( [ call( { "type": job_types.WORKFLOW_APPROVAL, "inputs": { "workflow_ingest_job_id": 88997690, }, }, session=mysql.db_session().__enter__(), ), call(jobs_to_persist, session=mysql.db_session().__enter__()), call( { "id": 13, "inputs": {"workflow_state_machine_execution_arn": "cool-arn"}, }, session=mysql.db_session().__enter__(), ), ] ) stepfunctions_input = { "workflow_ingest_job_id": 88997690, "workflow_job_type": job_types.WORKFLOW_APPROVAL, "workflow_job_id": parent_id, } start_execution_mock.assert_called_once_with( json.dumps(stepfunctions_input), "{}_{}".format("test_uuid", parent_id), job_types.WORKFLOW_APPROVAL, ) REENCODE_INGEST_JOB_ID = 55 REENCODE_PRIOR_VALIDATE_ANALYSIS_INPUTS = { job_io_fields.VIDEO_INNER_WIDTH_PIXELS: 1920, job_io_fields.VIDEO_INNER_HEIGHT_PIXELS: 800, job_io_fields.SILENT_INTERVAL_AT_BEGINNING_ENDS_AT_SECONDS: 0, job_io_fields.SILENT_INTERVAL_AT_END_BEGINS_AT_SECONDS: 100, job_io_fields.MAX_VOLUME_DECIBELS: -3, job_io_fields.BLACK_INTERVAL_AT_BEGINNING_ENDS_AT_SECONDS: 0, job_io_fields.BLACK_INTERVAL_AT_END_BEGINS_AT_SECONDS: 100, job_io_fields.LEFT_AUDIO_CHANNEL_RMS_VOLUME_DECIBELS: -18, job_io_fields.RIGHT_AUDIO_CHANNEL_RMS_VOLUME_DECIBELS: -18, job_io_fields.VIDEO_STREAM_FRAME_RATE_FRAMES_PER_SECOND: "24", job_io_fields.VIDEO_STREAM_PIXEL_ASPECT_RATIO: "1:1", } REENCODE_PRIOR_CREATE_PREVIEW_INPUTS = { job_io_fields.INPUT_VIDEO_S3_BUCKET: "bucket", job_io_fields.INPUT_VIDEO_S3_KEY: "raw/55/file.mov", job_io_fields.VIDEO_INNER_TOP_LEFT_CORNER_X_PIXELS: 0, job_io_fields.VIDEO_INNER_TOP_LEFT_CORNER_Y_PIXELS: 140, job_io_fields.VIDEO_INNER_WIDTH_PIXELS: 1920, job_io_fields.VIDEO_INNER_HEIGHT_PIXELS: 800, job_io_fields.VIDEO_STREAM_DURATION_SECONDS: 100, job_io_fields.VIDEO_STREAM_FRAME_RATE_FRAMES_PER_SECOND: "24", # The prior run's identity fields — must not leak into the new run's seed. job_io_fields.PRODUCT_ID: 888, job_io_fields.WORKFLOW_JOB_ID: 999, # validate_analysis re-derives these (including the deprecated variant); none # may survive into the seed. job_io_fields.DEPRECATED_VIDEO_OUTPUT_WIDTH_PIXELS: 1920, job_io_fields.VIDEO_OUTPUT_INNER_WIDTH_PIXELS: 1920, job_io_fields.VIDEO_OUTPUT_INNER_HEIGHT_PIXELS: 800, job_io_fields.VIDEO_OUTPUT_OUTER_WIDTH_PIXELS: 1920, job_io_fields.VIDEO_OUTPUT_OUTER_HEIGHT_PIXELS: 1080, job_io_fields.AUDIO_VOLUME_ADJUSTMENT_DB: 3, job_io_fields.INPUT_CLIP_START_TIMECODE: "00:00:00:00", job_io_fields.INPUT_CLIP_END_TIMECODE: "00:01:40:00", } # Crop-only overrides (timeline-preserving), so the chosen thumbnail re-points # to the new run. Black/silent (timeline-shifting) overrides are covered by the # _reencoded_thumbnail_path unit tests. REENCODE_OVERRIDES = { job_io_fields.VIDEO_INNER_TOP_LEFT_CORNER_Y_PIXELS: 0, job_io_fields.VIDEO_INNER_HEIGHT_PIXELS: 1080, } def _reencode_prior_jobs() -> list[dict[str, Any]]: return [ { "id": 500, "parent_id": REENCODE_INGEST_JOB_ID, "type": job_types.EXTRACT_METADATA, "status": job_statuses.COMPLETE, "inputs": {"should_not_be_seeded": True}, }, { "id": 501, "parent_id": REENCODE_INGEST_JOB_ID, "type": job_types.VALIDATE_ANALYSIS, "status": job_statuses.COMPLETE, "inputs": dict(REENCODE_PRIOR_VALIDATE_ANALYSIS_INPUTS), }, { "id": 502, "parent_id": REENCODE_INGEST_JOB_ID, "type": job_types.CREATE_STREAMABLE_PREVIEW_AND_THUMBNAILS, "status": job_statuses.COMPLETE, "inputs": dict(REENCODE_PRIOR_CREATE_PREVIEW_INPUTS), }, ] def test_setup_ingestion_reencode_workflow(mocker: MockerFixture) -> None: """The re-encode seeds validate_analysis + create_preview from a prior run. The overrides are applied (crop and interval alike), validate_analysis-derived fields are dropped so they recompute, the prior run's identity fields are replaced, unrelated prior jobs are ignored, the product's latest_pipeline_run_id is pointed at the new run, and the chosen thumbnail is repointed to that run so it shows the new crop. """ parent_id = 13 persist_mock = mocker.patch.object( queries, "persist_job_data", autospec=True, return_value=[parent_id] ) upsert_mock = mocker.patch.object( product_video, "upsert", autospec=True, return_value=parent_id ) mocker.patch.object( product_video, "get", autospec=True, return_value={"thumbnail_path": f"{REENCODE_INGEST_JOB_ID}/12.0000001.jpg"}, ) mocker.patch.object( queries, "get_jobs", autospec=True, side_effect=[ _reencode_prior_jobs(), [{"id": parent_id, "type": job_types.WORKFLOW_INGESTION_REENCODE}], ], ) mocker.patch.object(job_logic, "check_has_access", autospec=True) start_execution_mock = mocker.patch.object( stepfunctions, "start_execution", autospec=True, return_value={"executionArn": "cool-arn"}, ) mocker.patch.object(mysql, "db_session", autospec=True) mocker.patch.object(uuid, "uuid1", return_value="test_uuid", autospec=True) job_logic.setup_ingestion_reencode_workflow( { job_io_fields.WORKFLOW_INGEST_JOB_ID: REENCODE_INGEST_JOB_ID, "overrides": dict(REENCODE_OVERRIDES), "context": {"product_id": 123, "upc": 1234567890123}, } ) expected_seed = { job_io_fields.VIDEO_INNER_WIDTH_PIXELS: 1920, job_io_fields.VIDEO_INNER_HEIGHT_PIXELS: 1080, job_io_fields.SILENT_INTERVAL_AT_BEGINNING_ENDS_AT_SECONDS: 0, job_io_fields.SILENT_INTERVAL_AT_END_BEGINS_AT_SECONDS: 100, job_io_fields.MAX_VOLUME_DECIBELS: -3, job_io_fields.BLACK_INTERVAL_AT_BEGINNING_ENDS_AT_SECONDS: 0, job_io_fields.BLACK_INTERVAL_AT_END_BEGINS_AT_SECONDS: 100, job_io_fields.LEFT_AUDIO_CHANNEL_RMS_VOLUME_DECIBELS: -18, job_io_fields.RIGHT_AUDIO_CHANNEL_RMS_VOLUME_DECIBELS: -18, job_io_fields.VIDEO_STREAM_FRAME_RATE_FRAMES_PER_SECOND: "24", job_io_fields.VIDEO_STREAM_PIXEL_ASPECT_RATIO: "1:1", job_io_fields.INPUT_VIDEO_S3_BUCKET: "bucket", job_io_fields.INPUT_VIDEO_S3_KEY: "raw/55/file.mov", job_io_fields.VIDEO_INNER_TOP_LEFT_CORNER_X_PIXELS: 0, job_io_fields.VIDEO_INNER_TOP_LEFT_CORNER_Y_PIXELS: 0, job_io_fields.VIDEO_STREAM_DURATION_SECONDS: 100, job_io_fields.PRODUCT_ID: 123, job_io_fields.WORKFLOW_INGEST_JOB_ID: REENCODE_INGEST_JOB_ID, } persist_mock.assert_has_calls( [ call( { "type": job_types.WORKFLOW_INGESTION_REENCODE, "inputs": expected_seed, }, session=mysql.db_session().__enter__(), ), call( [ {"parent_id": parent_id, "type": job_types.VALIDATE_ANALYSIS}, { "parent_id": parent_id, "type": job_types.CREATE_STREAMABLE_PREVIEW_AND_THUMBNAILS, }, ], session=mysql.db_session().__enter__(), ), ] ) upsert_mock.assert_called_once_with( product_video={ "release_id": 123, "latest_pipeline_run_id": parent_id, "thumbnail_path": f"{parent_id}/12.0000001.jpg", } ) stepfunctions_input = { **expected_seed, "workflow_job_type": job_types.WORKFLOW_INGESTION_REENCODE, "workflow_job_id": parent_id, } start_execution_mock.assert_called_once_with( json.dumps(stepfunctions_input), "{}_{}".format("test_uuid", parent_id), job_types.WORKFLOW_INGESTION_REENCODE, ) @pytest.mark.parametrize( ("current_product", "overrides", "expected"), [ pytest.param( {"thumbnail_path": "55/12.0000001.jpg"}, {job_io_fields.VIDEO_INNER_HEIGHT_PIXELS: 1080}, {"thumbnail_path": "777/12.0000001.jpg"}, id="crop_only_repoints_frame_to_reencode_run", ), pytest.param( {"thumbnail_path": "55/custom.jpg"}, {job_io_fields.VIDEO_INNER_HEIGHT_PIXELS: 1080}, {}, id="custom_upload_left_untouched", ), pytest.param( {"thumbnail_path": "55/12.0000001.jpg"}, {job_io_fields.INPUT_CLIP_START_TIMECODE: "00:00:02:00"}, {"thumbnail_path": None, "thumbnail_at_milliseconds": None}, id="clip_override_clears_the_thumbnail", ), pytest.param( {"thumbnail_path": "55/12.0000001.jpg"}, { job_io_fields.VIDEO_INNER_HEIGHT_PIXELS: 1080, job_io_fields.INPUT_CLIP_END_TIMECODE: "00:01:38:00", }, {"thumbnail_path": None, "thumbnail_at_milliseconds": None}, id="crop_plus_clip_override_clears_the_thumbnail", ), pytest.param( {}, {job_io_fields.VIDEO_INNER_HEIGHT_PIXELS: 1080}, {}, id="no_chosen_thumbnail", ), ], ) def test_reencode_thumbnail_updates( mocker: MockerFixture, current_product: dict[str, Any], overrides: dict[str, Any], expected: dict[str, Any], ) -> None: mocker.patch.object( product_video, "get", autospec=True, return_value=current_product ) result = job_logic._reencode_thumbnail_updates( product_id=123, reencode_job_id=777, overrides=overrides ) assert result == expected def test_setup_ingestion_reencode_workflow_missing_prior_render( mocker: MockerFixture, ) -> None: """A referenced run without a create_preview job cannot be re-encoded.""" mocker.patch.object( queries, "get_jobs", autospec=True, return_value=[ { "id": 501, "parent_id": REENCODE_INGEST_JOB_ID, "type": job_types.VALIDATE_ANALYSIS, "status": job_statuses.COMPLETE, "inputs": dict(REENCODE_PRIOR_VALIDATE_ANALYSIS_INPUTS), }, ], ) with pytest.raises(InvalidRequest): job_logic.setup_ingestion_reencode_workflow( { job_io_fields.WORKFLOW_INGEST_JOB_ID: REENCODE_INGEST_JOB_ID, "overrides": dict(REENCODE_OVERRIDES), "context": {"product_id": 123}, } ) def test_setup_ingestion_reencode_workflow_incomplete_prior_run( mocker: MockerFixture, ) -> None: """A run whose render did not complete cannot be re-encoded from.""" prior_jobs = _reencode_prior_jobs() prior_jobs[-1]["status"] = job_statuses.ERROR mocker.patch.object(queries, "get_jobs", autospec=True, return_value=prior_jobs) with pytest.raises(InvalidRequest): job_logic.setup_ingestion_reencode_workflow( { job_io_fields.WORKFLOW_INGEST_JOB_ID: REENCODE_INGEST_JOB_ID, "overrides": dict(REENCODE_OVERRIDES), "context": {"product_id": 123}, } ) @pytest.mark.parametrize( ("test_description", "request_data"), [ ( "missing workflow_ingest_job_id", { "overrides": {job_io_fields.VIDEO_INNER_HEIGHT_PIXELS: 1080}, "context": {"product_id": 123}, }, ), ( "missing overrides", { job_io_fields.WORKFLOW_INGEST_JOB_ID: REENCODE_INGEST_JOB_ID, "context": {"product_id": 123}, }, ), ( "empty overrides", { job_io_fields.WORKFLOW_INGEST_JOB_ID: REENCODE_INGEST_JOB_ID, "overrides": {}, "context": {"product_id": 123}, }, ), ( "missing product_id", { job_io_fields.WORKFLOW_INGEST_JOB_ID: REENCODE_INGEST_JOB_ID, "overrides": {job_io_fields.VIDEO_INNER_HEIGHT_PIXELS: 1080}, "context": {}, }, ), ( "override of a non-threshold field", { job_io_fields.WORKFLOW_INGEST_JOB_ID: REENCODE_INGEST_JOB_ID, "overrides": {job_io_fields.MAX_VOLUME_DECIBELS: -2}, "context": {"product_id": 123}, }, ), ( "override of a raw black interval (replaced by the clip abstraction)", { job_io_fields.WORKFLOW_INGEST_JOB_ID: REENCODE_INGEST_JOB_ID, "overrides": { job_io_fields.BLACK_INTERVAL_AT_BEGINNING_ENDS_AT_SECONDS: 2, }, "context": {"product_id": 123}, }, ), ], ) def test_setup_ingestion_reencode_workflow_invalid_request( mocker: MockerFixture, test_description: str, request_data: dict[str, Any] ) -> None: """Malformed re-encode requests are rejected before any work is done.""" get_jobs_mock = mocker.patch.object(queries, "get_jobs", autospec=True) with pytest.raises(jsonschema.ValidationError): job_logic.setup_ingestion_reencode_workflow(request_data) get_jobs_mock.assert_not_called() def test_get_cloud_transfer_progress(mocker: MockerFixture) -> None: """Test get cloud transfer progress.""" mocker.patch.object( cloud_transfer_progress, "get_cloud_transfer_progress", return_value="testing", autospec=True, ) assert "testing" == job_logic.get_cloud_transfer_progress(123, "asdf")