"""Unit tests for src/app.py.""" import json from httpx import HTTPStatusError, Response from unittest.mock import MagicMock, patch import pytest from src import app def make_mock_message(asset_final_id=12345, filename="test.jpg"): """Create a mock Kafka MSK message.""" message = MagicMock() message.value = json.dumps( { "ASSET_FINAL_ID": asset_final_id, "FILENAME": filename, "ASSET_TYPE": "some-type-0", "ASSET_SUBTYPE": "some-subtype-0", } ) message.topic = "test-topic" return message def make_event_payload(asset_final_id=12345, filename="test.jpg"): """Create an AssetMezzanineCoverartEventMessage payload for process_event tests.""" return app.AssetMezzanineCoverartEventMessage.from_dict( { "ASSET_FINAL_ID": asset_final_id, "FILENAME": filename, "ASSET_TYPE": "some-type-0", "ASSET_SUBTYPE": "some-subtype-0", } ) # --------------------------------------------------------------------------- # process_event # --------------------------------------------------------------------------- @patch("src.app.lambda_logger") @patch("src.app.ows_assets") @patch("src.app.hive") def test_process_event_success(mock_hive, mock_ows_assets, mock_logger): """process_event fetches a presigned URL, runs the Hive task, and posts results.""" asset_final_id = 12345 signed_url = "http://s3.example.com/asset.jpg" block_text = "Sample extracted text" payload = make_event_payload(asset_final_id) mock_ows_assets.get_presigned_url.return_value = signed_url mock_hive.run_task.return_value = block_text app.process_event(payload) mock_ows_assets.get_presigned_url.assert_called_once_with(asset_final_id) mock_hive.run_task.assert_called_once_with(signed_url) mock_ows_assets.post_results.assert_called_once_with(asset_final_id, block_text) mock_logger.set_data.assert_called_with( status="success", result="Saved Hive result to ows-assets." ) @patch("src.app.lambda_logger") @patch("src.app.ows_assets") @patch("src.app.hive") def test_process_event_skip_when_result_already_exists( mock_hive, mock_ows_assets, mock_logger ): """process_event logs a skip when post_results returns 409 Conflict.""" asset_final_id = 12345 signed_url = "http://s3.example.com/asset.jpg" block_text = "Sample extracted text" payload = make_event_payload(asset_final_id) mock_ows_assets.get_presigned_url.return_value = signed_url mock_hive.run_task.return_value = block_text mock_ows_assets.post_results.side_effect = HTTPStatusError( "Conflict", request=MagicMock(), response=Response(409) ) app.process_event(payload) mock_ows_assets.post_results.assert_called_once_with(asset_final_id, block_text) mock_logger.set_data.assert_called_with( status="skip", result="Existing Hive result found for asset." ) @patch("src.app.lambda_logger") @patch("src.app.ows_assets") @patch("src.app.hive") def test_process_event_raises_on_non_409_post_error( mock_hive, mock_ows_assets, mock_logger ): """process_event re-raises HTTP errors other than 409 from post_results.""" asset_final_id = 12345 signed_url = "http://s3.example.com/asset.jpg" block_text = "Sample extracted text" payload = make_event_payload(asset_final_id) mock_ows_assets.get_presigned_url.return_value = signed_url mock_hive.run_task.return_value = block_text mock_ows_assets.post_results.side_effect = HTTPStatusError( "Internal Server Error", request=MagicMock(), response=Response(500) ) with pytest.raises(HTTPStatusError): app.process_event(payload) @patch("src.app.lambda_logger") @patch("src.app.ows_assets") @patch("src.app.hive") def test_process_event_raises_on_get_presigned_url_error( mock_hive, mock_ows_assets, mock_logger ): """process_event re-raises non-404 HTTP errors from get_presigned_url.""" asset_final_id = 12345 payload = make_event_payload(asset_final_id) mock_ows_assets.get_presigned_url.side_effect = HTTPStatusError( "Internal Server Error", request=MagicMock(), response=Response(500) ) with pytest.raises(HTTPStatusError): app.process_event(payload) mock_hive.run_task.assert_not_called() mock_ows_assets.post_results.assert_not_called() @patch("src.app.lambda_logger") @patch("src.app.ows_assets") @patch("src.app.hive") def test_process_event_skip_when_asset_not_found_non_prod( mock_hive, mock_ows_assets, mock_logger ): """process_event logs a skip when get_presigned_url returns 404 in non-prod.""" payload = make_event_payload(12345) mock_ows_assets.get_presigned_url.side_effect = HTTPStatusError( "Not Found", request=MagicMock(), response=Response(404) ) with patch("src.app.config.ENVIRONMENT", "test"): app.process_event(payload) mock_hive.run_task.assert_not_called() mock_ows_assets.post_results.assert_not_called() mock_logger.set_data.assert_called_with( status="skip", result="Failed to get download url: asset not found." ) @patch("src.app.lambda_logger") @patch("src.app.ows_assets") @patch("src.app.hive") def test_process_event_raises_on_404_in_prod(mock_hive, mock_ows_assets, mock_logger): """process_event re-raises 404 from get_presigned_url when in prod environment.""" payload = make_event_payload(12345) mock_ows_assets.get_presigned_url.side_effect = HTTPStatusError( "Not Found", request=MagicMock(), response=Response(404) ) with patch("src.app.config.ENVIRONMENT", "prod"): with pytest.raises(HTTPStatusError): app.process_event(payload) mock_hive.run_task.assert_not_called() mock_ows_assets.post_results.assert_not_called() # --------------------------------------------------------------------------- # handler # --------------------------------------------------------------------------- @patch("src.app.EventSourceMessage") @patch("src.app.lambda_logger") @patch("src.app.ows_assets") @patch("src.app.hive") def test_handler_success( mock_hive, mock_ows_assets, mock_logger, mock_event_source, ): """handler processes messages and returns STATUS_OK on success.""" asset_final_id = 12345 signed_url = "http://s3.example.com/asset.jpg" block_text = "Sample extracted text" message = make_mock_message(asset_final_id) event_key = "test-topic-0" mock_event_source.return_value.__iter__ = MagicMock( return_value=iter([(event_key, message)]) ) mock_ows_assets.get_presigned_url.return_value = signed_url mock_hive.run_task.return_value = block_text result = app.handler({"eventSource": "aws:kafka"}, None) assert result == {"status": "ok"} mock_ows_assets.get_presigned_url.assert_called_once_with(asset_final_id) mock_hive.run_task.assert_called_once_with(signed_url) mock_ows_assets.post_results.assert_called_once_with(asset_final_id, block_text) @patch("src.app.EventSourceMessage") @patch("src.app.lambda_logger") @patch("src.app.ows_assets") @patch("src.app.hive") def test_handler_empty_event( mock_hive, mock_ows_assets, mock_logger, mock_event_source, ): """handler returns STATUS_OK when the event contains no records.""" mock_event_source.return_value.__iter__ = MagicMock(return_value=iter([])) result = app.handler({"eventSource": "aws:kafka", "records": {}}, None) assert result == {"status": "ok"} mock_ows_assets.get_presigned_url.assert_not_called() mock_hive.run_task.assert_not_called() mock_ows_assets.post_results.assert_not_called() @patch("src.app.EventSourceMessage") @patch("src.app.lambda_logger") @patch("src.app.ows_assets") @patch("src.app.hive") def test_handler_logs_error_and_raises_on_exception( mock_hive, mock_ows_assets, mock_logger, mock_event_source, ): """handler sets STATUS_ERROR on the logger and re-raises unexpected exceptions.""" message = make_mock_message(12345) event_key = "test-topic-0" mock_event_source.return_value.__iter__ = MagicMock( return_value=iter([(event_key, message)]) ) mock_ows_assets.get_presigned_url.side_effect = HTTPStatusError( "Internal Server Error", request=MagicMock(), response=Response(500) ) with pytest.raises(HTTPStatusError): app.handler({"eventSource": "aws:kafka"}, None) mock_logger.set_data.assert_called_with( status="error", result=mock_ows_assets.get_presigned_url.side_effect.args[0] ) mock_logger.end.assert_called_once() @patch("src.app.EventSourceMessage") @patch("src.app.lambda_logger") @patch("src.app.ows_assets") @patch("src.app.hive") def test_handler_calls_logger_end_on_success( mock_hive, mock_ows_assets, mock_logger, mock_event_source, ): """handler always calls lambda_logger.end() via the finally block.""" message = make_mock_message(12345) event_key = "test-topic-0" mock_event_source.return_value.__iter__ = MagicMock( return_value=iter([(event_key, message)]) ) mock_ows_assets.get_presigned_url.return_value = "http://s3.example.com/asset.jpg" mock_hive.run_task.return_value = "some text" app.handler({"eventSource": "aws:kafka"}, None) mock_logger.end.assert_called_once() @patch("src.app.process_event") @patch("src.app.AssetMezzanineCoverartEventMessage.from_dict") @patch("src.app.EventSourceMessage") def test_handler_custom_event_processes_assets( mock_event_source, mock_from_dict, mock_process_event, ): """handler processes custom assets and does not use EventSourceMessage.""" asset_1 = {"ASSET_FINAL_ID": 1001} asset_2 = {"ASSET_FINAL_ID": 1002} data_1 = MagicMock(asset_final_id=1001) data_2 = MagicMock(asset_final_id=1002) mock_from_dict.side_effect = [data_1, data_2] result = app.handler({"eventSource": "custom", "assets": [asset_1, asset_2]}, None) assert result == {"status": "ok"} mock_event_source.assert_not_called() assert mock_from_dict.call_count == 2 mock_from_dict.assert_any_call(asset_1) mock_from_dict.assert_any_call(asset_2) mock_process_event.assert_any_call(data_1) mock_process_event.assert_any_call(data_2) @patch("src.app.process_event") @patch("src.app.AssetMezzanineCoverartEventMessage.from_dict") def test_handler_custom_event_raises_on_process_error( mock_from_dict, mock_process_event, ): """handler re-raises exceptions raised while processing custom assets.""" asset = {"ASSET_FINAL_ID": 1001} data = MagicMock(asset_final_id=1001) mock_from_dict.return_value = data mock_process_event.side_effect = RuntimeError("boom") with pytest.raises(RuntimeError, match="boom"): app.handler({"eventSource": "custom", "assets": [asset]}, None) mock_from_dict.assert_called_once_with(asset) mock_process_event.assert_called_once_with(data)