"""Test tasks for Deezer Snowflake-only Ingestion Workflow.""" from datetime import date as date_module, timedelta from importlib import reload import os import shutil from shutil import ReadError as ZipReadError from tempfile import gettempdir from unittest.mock import MagicMock from unittest.mock import Mock from unittest.mock import patch from garcon_contrib.dynamo_feed_status import garcon_feed_status import pytest from feed_ingestion import tasks as common_tasks from feed_ingestion.flows.deezer import tasks from feed_ingestion.tasks import deezer_tasks, s3_tasks from feed_ingestion.util import os_tools from feed_ingestion.util import task_status from feed_ingestion.util.aws import s3 as s3utils import feed_ingestion.util.deezer_zephir_utils as utils @pytest.fixture(autouse=True) def mock_check_status(request, monkeypatch): """Mock check_status decorator. Mock so that the decorated function is kept intact and always get called. """ mock = MagicMock() # just return the function without any modifications mock.return_value = lambda f: f monkeypatch.setattr(common_tasks, 'check_status', value=mock) reload(tasks) # redecorate tasks yield monkeypatch.undo() reload(tasks) @pytest.fixture def mock_boto3(): """Mock boto3.""" boto3_path = ( 'feed_ingestion.flows.amazon_digital_services.tasks.boto3') with patch(boto3_path) as boto3: mock_client = MagicMock() mock_client.send_email = MagicMock() boto3.client.return_value = mock_client yield boto3 @pytest.fixture def mock_set_overall_status(): """Yield overall status.""" overall_status_path = ( 'feed_ingestion.flows.deezer.tasks.garcon_feed_status.' 'set_overall_status') with patch(overall_status_path) as overall_status: yield overall_status @pytest.fixture def mock_set_missing_files(): """Yield set_missing_files.""" path = ( 'feed_ingestion.flows.deezer.tasks.garcon_feed_status.' 'set_missing_files') with patch(path) as set_missing_files: yield set_missing_files @pytest.fixture(params=[ ('2015-02-20', '2015-02-20'), ('2015-11-20', '2015-11-20'), ('', (date_module.today() - timedelta(days=1)).strftime('%Y-%m-%d')) ]) def context_date(request): """Fixture returning different dates. Fixture returning different dates as input and corresponding expected dates coming back from the bootstrap task. """ return request.param @patch('feed_ingestion.flows.deezer.tasks.garcon_feed_status') def test_bootstrap_theorchard(garcon_feed_status_mock, context_date): """Test bootstrap task.""" activity_mock = MagicMock() input_date, expected_date = ('2017-01-01', '2017-01-01') # context_date response = tasks.bootstrap(activity_mock, input_date, 'theorchard', None) assert 'feed_name' in response assert 'secrets_path' in response assert 'source_path' in response assert expected_date in response['s3_archive_bucket'] assert expected_date in response['s3_temp_staging_raw_bucket'] assert expected_date.replace('-', '') in response['drop_file_name'] assert response['spec_version'] == 1 assert response['date'] == expected_date assert response['staging_raw_table'] == 'staging_raw_deezer_v1' @patch('feed_ingestion.flows.deezer.tasks.garcon_feed_status') def test_bootstrap_sme(garcon_feed_status_mock, context_date): """Test bootstrap task.""" activity_mock = MagicMock() input_date, expected_date = ('2017-01-01', '2017-01-01') # context_date response = tasks.bootstrap(activity_mock, input_date, 'sme', None) assert 'feed_name' in response assert 'secrets_path' in response assert 'source_path' in response assert expected_date in response['s3_archive_bucket'] assert expected_date in response['s3_temp_staging_raw_bucket'] assert expected_date.replace('-', '') in response['drop_file_name'] assert response['spec_version'] == 2 assert response['date'] == expected_date assert response['staging_raw_table'] == 'staging_raw_deezer_v2' @patch('feed_ingestion.flows.deezer.tasks.garcon_feed_status') def test_bootstrap_altafonte(garcon_feed_status_mock): """Test bootstrap task.""" activity_mock = MagicMock() input_date, expected_date = ('2024-06-24', '2024-06-24') response = tasks.bootstrap(activity_mock, input_date, 'altafonte', None) assert 'feed_name' in response assert 'secrets_path' in response assert response[ 'source_path'] == 'ftp/altafonte/merlin/dzr-deezer/2024-06-24/' assert expected_date in response['s3_archive_bucket'] assert expected_date in response['s3_temp_staging_raw_bucket'] assert expected_date.replace('-', '') in response['drop_file_name'] assert response['spec_version'] == 4 assert response['date'] == expected_date assert response['staging_raw_table'] == 'staging_raw_deezer_v4' assert response['source_files_dict']['files'] == [ 'AltafonteMERLIN_20240624_20240624.txt', 'AltafonteMERLIN_20240624_20240624_TB.txt'] @patch.object(tasks, 'os') def test_fetch_from_drop_location(os_mock, monkeypatch, context_date): """Test fetch_from_drop_location: normal flow.""" # setup mocks/fixtures activity_mock = MagicMock() monkeypatch.setattr( garcon_feed_status, 'get_overall_status', value=MagicMock()) input_date, expected_date = context_date context = tasks.bootstrap( activity_mock, input_date, 'theorchard') create_temp_dir = MagicMock(return_value=os.path.join( gettempdir(), str(context.get('feed_name')), str(context.get('date')))) # mock download_daily_report_from_zephir response download_daily_report_from_zephir = MagicMock() download_daily_report_from_zephir.return_value = os.path.join( create_temp_dir(), str(context.get('drop_file_name'))) upload_to_s3 = MagicMock(return_value='16157363') delete_s3_obj = MagicMock(return_value='file deleted') # execute monkeypatch.setattr(task_status, 'is_completed_task', MagicMock( return_value=False)) monkeypatch.setattr( garcon_feed_status, 'set_overall_status', value=MagicMock()) monkeypatch.setattr(os_tools, 'create_temp_dir', create_temp_dir) monkeypatch.setattr( utils, 'download_daily_report_from_zephir', download_daily_report_from_zephir) monkeypatch.setattr( s3utils, 'upload_to_s3', upload_to_s3) monkeypatch.setattr( s3utils, 'delete_s3_obj', delete_s3_obj) monkeypatch.setattr(task_status, 'mark_completed_task', value=MagicMock()) tasks.fetch_from_drop_location( activity_mock, context.get('date'), context.get('feed_name'), context.get('source_path'), context.get('s3_archive_bucket'), context.get('drop_file_name') ) assert task_status.is_completed_task.called assert delete_s3_obj.called assert os_tools.create_temp_dir.called assert download_daily_report_from_zephir.called assert upload_to_s3.called assert os_mock.remove.called assert task_status.mark_completed_task.called def test_fetch_from_drop_location_already_complete(monkeypatch, context_date): """Test fetch_from_drop_location: task already completed.""" # setup mocks/fixtures activity_mock = MagicMock() monkeypatch.setattr( garcon_feed_status, 'get_overall_status', value=MagicMock()) input_date, expected_date = context_date context = tasks.bootstrap( activity_mock, input_date, 'theorchard') # execute monkeypatch.setattr(task_status, 'is_completed_task', MagicMock( return_value=True)) monkeypatch.setattr(task_status, 'mark_completed_task', value=MagicMock()) monkeypatch.setattr(utils, 'download_daily_report_from_zephir', value=MagicMock()) tasks.fetch_from_drop_location( activity_mock, context.get('date'), context.get('feed_name'), context.get('source_path'), context.get('s3_archive_bucket'), context.get('drop_file_name') ) assert not utils.download_daily_report_from_zephir.called assert not task_status.mark_completed_task.called def test_fetch_from_drop_location_download_failure(monkeypatch, context_date): """Test fetch_from_drop_location: zephir fetch failure.""" # setup mocks/fixtures activity_mock = MagicMock() monkeypatch.setattr( garcon_feed_status, 'get_overall_status', value=MagicMock()) input_date, expected_date = context_date context = tasks.bootstrap( activity_mock, input_date, 'theorchard') upload_to_s3 = MagicMock(return_value='16157363') delete_s3_obj = MagicMock(return_value='file deleted') # mock download_daily_report_from_zephir response download_daily_report_from_zephir = Mock(side_effect=FileNotFoundError()) monkeypatch.setattr(task_status, 'is_completed_task', MagicMock( return_value=False)) monkeypatch.setattr( garcon_feed_status, 'set_missing_files', value=MagicMock()) monkeypatch.setattr( garcon_feed_status, 'set_overall_status', value=MagicMock()) monkeypatch.setattr(task_status, 'mark_completed_task', value=MagicMock()) monkeypatch.setattr( utils, 'download_daily_report_from_zephir', download_daily_report_from_zephir) monkeypatch.setattr( s3utils, 'upload_to_s3', upload_to_s3) monkeypatch.setattr( s3utils, 'delete_s3_obj', delete_s3_obj) tasks.fetch_from_drop_location( activity_mock, context.get('date'), context.get('feed_name'), context.get('source_path'), context.get('s3_archive_bucket'), context.get('drop_file_name') ) assert task_status.is_completed_task.called assert download_daily_report_from_zephir.called garcon_feed_status.set_overall_status.assert_called_with( context.get('feed_name'), expected_date, garcon_feed_status.STATUS_NOT_AVAILABLE) assert not task_status.mark_completed_task.called def test_flow_ingested(monkeypatch): """Test stop message returned if flow already ingested.""" activity_mock = MagicMock() monkeypatch.setattr( garcon_feed_status, 'get_overall_status', value=MagicMock(return_value=garcon_feed_status.STATUS_INGESTED)) context = tasks.fetch_from_drop_location( activity_mock, 'date', 'feed_name', 'source_path', 'target_s3_path', 'drop_file_name') assert context['stop'] def test_grab_drop_files_sme_files_exist( monkeypatch, context_date, mock_set_missing_files, mock_set_overall_status): """Test grab_drop_files_sme: files exist.""" activity_mock = MagicMock() monkeypatch.setattr( garcon_feed_status, 'get_overall_status', value=MagicMock()) input_date, expected_date = context_date context = tasks.bootstrap(activity_mock, input_date, 'theorchard') extract_bucket_path_mock = Mock() extract_bucket_path_mock.return_value = ( 'key_name', 'bucket_name') copy_file_from_sme_s3_to_theocrhard_mock = Mock() copy_file_from_sme_s3_to_theocrhard_mock.return_value = \ {context.get('source_files_dict').get('files')[0]: True} monkeypatch.setattr( s3_tasks, 'copy_file_from_sme_s3_to_theocrhard', copy_file_from_sme_s3_to_theocrhard_mock) result = tasks.grab_drop_files_sme( activity_mock, context.get('feed_name'), context.get('date'), context.get('source_files_dict').get('files')[0], context.get('s3_archive_bucket'), context.get('source_path'), context.get('s3_archive_bucket') ) assert copy_file_from_sme_s3_to_theocrhard_mock.called mock_set_missing_files.assert_not_called() assert result is None def test_grab_drop_files_sme_files_do_not_exist( monkeypatch, context_date, mock_set_missing_files, mock_set_overall_status): """Test grab_drop_files_sme: files doesn't exist.""" activity_mock = MagicMock() monkeypatch.setattr( garcon_feed_status, 'get_overall_status', value=MagicMock()) input_date, expected_date = context_date context = tasks.bootstrap(activity_mock, input_date, 'theorchard') extract_bucket_path_mock = Mock() extract_bucket_path_mock.return_value = ( 'key_name', 'bucket_name') copy_file_from_sme_s3_to_theocrhard_mock = Mock() copy_file_from_sme_s3_to_theocrhard_mock.return_value = \ {context.get('source_files_dict').get('files')[0]: False} monkeypatch.setattr( s3_tasks, 'copy_file_from_sme_s3_to_theocrhard', copy_file_from_sme_s3_to_theocrhard_mock) result = tasks.grab_drop_files_sme( activity_mock, context.get('feed_name'), context.get('date'), context.get('source_files_dict').get('files')[0], context.get('s3_archive_bucket'), context.get('source_path'), context.get('s3_archive_bucket') ) assert copy_file_from_sme_s3_to_theocrhard_mock.called garcon_feed_status.set_overall_status.assert_called_with( context.get('feed_name'), expected_date, garcon_feed_status.STATUS_NOT_AVAILABLE) mock_set_missing_files.assert_called() assert result == {'stop': True} @patch('feed_ingestion.flows.deezer.tasks.garcon_feed_status') @patch.object(tasks, 'os') def test_unzip_and_clean( os_mock, garcon_feed_status_mock, monkeypatch, context_date ): """Test unzip_and_clean: successful operation.""" activity_mock = MagicMock() input_date, expected_date = context_date context = tasks.bootstrap(activity_mock, input_date, 'theorchard') create_temp_dir = MagicMock(return_value=os.path.join( gettempdir(), str(context.get('feed_name')), str(context.get('date')))) upload_to_s3 = MagicMock(return_value='16157363') download_from_s3 = MagicMock(return_value='16157363') delete_s3_obj = MagicMock(return_value='file deleted') files = [ { 'source_file_name': 'source_file_name', 'gzip_file_name': 'gzip_file_name', 'gzip_file': 'gzip_file' }, { 'source_file_name': 'source_file_name', 'gzip_file_name': 'gzip_file_name', 'gzip_file': 'gzip_file' } ] monkeypatch.setattr(os_tools, 'create_temp_dir', create_temp_dir) monkeypatch.setattr( s3utils, 'download_from_s3', download_from_s3) monkeypatch.setattr( s3utils, 'upload_to_s3', upload_to_s3) monkeypatch.setattr( s3utils, 'delete_s3_obj', delete_s3_obj) monkeypatch.setattr( utils, 'repack_source_file', MagicMock(return_value=files)) monkeypatch.setattr( shutil, 'rmtree', MagicMock(return_value='')) monkeypatch.setattr( garcon_feed_status, 'set_overall_status', value=MagicMock()) deezer_tasks.unzip_and_clean( activity_mock, context.get('date'), context.get('feed_name'), context.get('s3_archive_bucket'), context.get('s3_temp_staging_raw_bucket'), context.get('drop_file_name'), context.get('licensor'), context.get('source_files_dict'), ) assert os_tools.create_temp_dir.call_count == 1 assert download_from_s3.call_count == 1 assert utils.repack_source_file.call_count == 1 assert delete_s3_obj.call_count == 2 assert upload_to_s3.call_count == 2 assert shutil.rmtree.call_count == 1 @patch('feed_ingestion.flows.deezer.tasks.garcon_feed_status') @patch.object(tasks, 'os') def test_unzip_and_clean_source_missing( os_mock, garcon_feed_status_mock, monkeypatch, context_date ): """Test unzip_and_clean: S3 source missing.""" activity_mock = MagicMock() input_date, expected_date = context_date context = tasks.bootstrap(activity_mock, input_date, 'theorchard') create_temp_dir = MagicMock(return_value=os.path.join( gettempdir(), str(context.get('feed_name')), str(context.get('date')))) upload_to_s3 = MagicMock(return_value='16157363') download_from_s3 = MagicMock(return_value=None) delete_s3_obj = MagicMock(return_value='file deleted') monkeypatch.setattr(os_tools, 'create_temp_dir', create_temp_dir) monkeypatch.setattr( s3utils, 'download_from_s3', download_from_s3) monkeypatch.setattr( s3utils, 'upload_to_s3', upload_to_s3) monkeypatch.setattr( s3utils, 'delete_s3_obj', delete_s3_obj) monkeypatch.setattr( utils, 'repack_source_file', MagicMock(return_value=[])) monkeypatch.setattr( garcon_feed_status, 'set_overall_status', value=MagicMock()) monkeypatch.setattr( shutil, 'rmtree', MagicMock(return_value='')) response = {} context['src_s3_file_key'] = 'file_name' response = deezer_tasks.unzip_and_clean( activity_mock, context.get('date'), context.get('feed_name'), context.get('s3_archive_bucket'), context.get('s3_temp_staging_raw_bucket'), context.get('drop_file_name'), context.get('licensor'), context.get('source_files_dict'), ) details = "'{}' download from '{}' bucket has failed".format( context.get('drop_file_name'), context.get('s3_archive_bucket')) assert os_tools.create_temp_dir.call_count == 1 assert download_from_s3.call_count == 1 assert utils.repack_source_file.call_count == 0 assert delete_s3_obj.call_count == 0 assert upload_to_s3.call_count == 0 assert shutil.rmtree.call_count == 0 assert response == {'message': details, 'stop': True} @patch('feed_ingestion.flows.deezer.tasks.garcon_feed_status') def test_unzip_and_clean_unzip_fails( garcon_feed_status_mock, monkeypatch, context_date): """Test unzip_and_clean: unzip operation fails.""" activity_mock = MagicMock() input_date, expected_date = context_date context = tasks.bootstrap(activity_mock, input_date, 'theorchard') create_temp_dir = MagicMock(return_value=os.path.join( gettempdir(), str(context.get('feed_name')), str(context.get('date')))) upload_to_s3 = MagicMock(return_value='16157363') download_from_s3 = MagicMock(return_value='1615736') delete_s3_obj = MagicMock(return_value='file deleted') monkeypatch.setattr(os_tools, 'create_temp_dir', create_temp_dir) monkeypatch.setattr( s3utils, 'download_from_s3', download_from_s3) monkeypatch.setattr( s3utils, 'upload_to_s3', upload_to_s3) monkeypatch.setattr( s3utils, 'delete_s3_obj', delete_s3_obj) monkeypatch.setattr( utils, 'repack_source_file', MagicMock(return_value=[])) monkeypatch.setattr( shutil, 'rmtree', MagicMock(return_value='')) monkeypatch.setattr( garcon_feed_status, 'set_overall_status', value=MagicMock()) response = {} try: response = deezer_tasks.unzip_and_clean( activity_mock, context.get('date'), context.get('feed_name'), context.get('s3_archive_bucket'), context.get('s3_temp_staging_raw_bucket'), context.get('drop_file_name'), context.get('licensor'), context.get('source_files_dict'), ) except ZipReadError as e: assert str(e) == 'File is not a zip file' pass assert os_tools.create_temp_dir.call_count == 1 assert download_from_s3.call_count == 1 assert utils.repack_source_file.call_count == 1 assert delete_s3_obj.call_count == 0 assert upload_to_s3.call_count == 0 assert shutil.rmtree.call_count == 1 assert response == {'source_files_dict': {'files': []}} @pytest.mark.parametrize('licensor,expected_fraud_name', [ ('sme', 'fraud_report_sony-20260504.txt'), ('theorchard', 'fraud_report_theorchard-20260504.txt'), ]) @patch('feed_ingestion.flows.deezer.tasks.garcon_feed_status') def test_unzip_and_clean_fraud_report_name_expansion( garcon_feed_status_mock, licensor, expected_fraud_name, monkeypatch ): """Test fraud_report_*.txt wildcard expansion. Uses Deezer's licensor identifier, not the internal licensor name (e.g. sony for sme). """ activity_mock = MagicMock() date = '2026-05-04' source_files_dict = { 'files': [ 'sony_activity_metrics_daily.csv', 'sony_detailed_content_identifiers_daily.csv', 'sony_playlist_metrics.csv', 'sony_stream_source.csv', 'sony_user_activity_metrics.csv', 'sony_user_details.csv', 'sony_vendor_product_details_daily.csv', 'fraud_report_*.txt', ] } repack_source_file_mock = MagicMock(return_value=[]) download_from_s3_mock = MagicMock(return_value='123') upload_to_s3_mock = MagicMock(return_value='123') delete_s3_obj_mock = MagicMock() monkeypatch.setattr( os_tools, 'create_temp_dir', MagicMock(return_value='/tmp/test')) monkeypatch.setattr(s3utils, 'download_from_s3', download_from_s3_mock) monkeypatch.setattr(s3utils, 'upload_to_s3', upload_to_s3_mock) monkeypatch.setattr(s3utils, 'delete_s3_obj', delete_s3_obj_mock) monkeypatch.setattr(utils, 'repack_source_file', repack_source_file_mock) monkeypatch.setattr(shutil, 'rmtree', MagicMock()) deezer_tasks.unzip_and_clean( activity_mock, date, f'deezer_daily_{licensor}', 's3://cucumbers/DeezerV3/archives/2026-05-04/', 's3://cucumbers/DeezerV3/temp/2026-05-04/', 'deezer_daily_report_activity_metrics_20260504.zip', licensor, source_files_dict, ) assert repack_source_file_mock.called actual_source = repack_source_file_mock.call_args[0][2] assert actual_source['files'][-1] == expected_fraud_name @patch('feed_ingestion.flows.deezer.tasks.garcon_feed_status') def test_bootstrap_fraud_report_backfill(garcon_feed_status_mock): """Test bootstrap with fraud_report_backfill flag.""" activity_mock = MagicMock() input_date = '2026-01-02' response = tasks.bootstrap( activity_mock, input_date, 'sme', None, fraud_report_backfill='True') assert response['fraud_report_backfill'] is True assert response['spec_version'] == 3 assert response['source_files_dict'] == {'files': ['fraud_report_*.txt']} assert response['staging_raw_table'] == [ 'staging_raw_deezer_v3_fraud_report'] assert response['feed_name'] == 'deezer_daily_sme' @pytest.mark.parametrize('licensor,expected_source_key,expected_target_gz', [ ( 'sme', 'deezer/in/sme/fraudulent_reports/20260102/' '20260102_fraud_report_sony-20260102.txt', 'fraud_report_sony-20260102.txt.gz', ), ( 'theorchard', 'deezer/in/orchard/fraudulent_reports/20260102/' '20260102_fraud_report_theorchard-20260102.txt', 'fraud_report_theorchard-20260102.txt.gz', ), ]) def test_grab_fraud_report_backfill_success( monkeypatch, mock_set_overall_status, mock_set_missing_files, licensor, expected_source_key, expected_target_gz): """Test grab_fraud_report_backfill downloads, gzips and uploads.""" activity_mock = MagicMock() sme_client_mock = MagicMock() monkeypatch.setattr( s3_tasks, '_get_sme_s3_client', Mock(return_value=sme_client_mock)) monkeypatch.setattr( os_tools, 'create_temp_dir', MagicMock(return_value='/tmp/test')) monkeypatch.setattr(s3utils, 'delete_s3_obj', MagicMock()) monkeypatch.setattr(s3utils, 'upload_to_s3', MagicMock(return_value='12345')) open_mock = MagicMock() monkeypatch.setattr('builtins.open', open_mock) monkeypatch.setattr('feed_ingestion.flows.deezer.tasks.gzip.open', MagicMock()) monkeypatch.setattr('feed_ingestion.flows.deezer.tasks.shutil.copyfileobj', MagicMock()) monkeypatch.setattr('feed_ingestion.flows.deezer.tasks.os.remove', MagicMock()) result = tasks.grab_fraud_report_backfill( activity_mock, f'deezer_daily_{licensor}', '2026-01-02', 'sme-ca-prod-partners', 's3://cucumbers/DeezerV3/temp/2026-01-02/sme/', licensor, ) assert result is None sme_client_mock.download_file.assert_called_once_with( 'sme-ca-prod-partners', expected_source_key, '/tmp/test/' + expected_source_key.split('/')[-1]) s3utils.upload_to_s3.assert_called_once() upload_args = s3utils.upload_to_s3.call_args[0] assert upload_args[0] == f'/tmp/test/{expected_target_gz}' assert upload_args[1] == 'cucumbers' assert upload_args[2].endswith(expected_target_gz) mock_set_overall_status.assert_not_called() def test_grab_fraud_report_backfill_file_missing( monkeypatch, mock_set_overall_status, mock_set_missing_files): """Test grab_fraud_report_backfill when source file is missing.""" activity_mock = MagicMock() sme_client_mock = MagicMock() sme_client_mock.download_file.side_effect = Exception('Not found') monkeypatch.setattr( s3_tasks, '_get_sme_s3_client', Mock(return_value=sme_client_mock)) monkeypatch.setattr( os_tools, 'create_temp_dir', MagicMock(return_value='/tmp/test')) result = tasks.grab_fraud_report_backfill( activity_mock, 'deezer_daily_sme', '2026-01-02', 'sme-ca-prod-partners', 's3://cucumbers/DeezerV3/temp/2026-01-02/sme/', 'sme', ) assert result == {'stop': True} mock_set_overall_status.assert_called_once() mock_set_missing_files.assert_called_once()