"""Unit tests for Pandora Analytics Ingestion Workflow.""" from unittest.mock import MagicMock from unittest.mock import patch from garcon.activity import Activity from garcon_contrib.dynamo_feed_status import garcon_feed_status from feed_ingestion.flows.pandora.flow import Flow _CONTEXT = { 'bootstrap.date': '2000-01-01', 'bootstrap.drop_bucket': 's3://drop_bucket/data/2000/01/01/', 'bootstrap.snowflake_error_on_column_count_mismatch': 'true', 'bootstrap.archive_bucket': 's3://data_bucket/Pandora/archives/2000-01-01/', 'bootstrap.expected_files': [ 'orchard_metadata_2000-01-01.txt.bz2', 'orchard_US_2000-01-01.txt.bz2'], 'bootstrap.temp_staging_raw_tables': { 'metadata': { 'temp_table_name': 'pandora_metadata_20000101', 'temp_table_s3_full_path': 's3://dev-cucumbers/Pandora/archives/2000-01-01/' 'orchard_metadata_2000-01-01.txt.bz2', 'kwargs': {'error_on_column_count_mismatch': 'true'} }, 'US': { 'temp_table_name': 'pandora_streams_US_20000101', 'temp_table_s3_full_path': 's3://dev-cucumbers/Pandora/archives/2000-01-01/' 'orchard_US_2000-01-01.txt.bz2', 'kwargs': {'error_on_column_count_mismatch': 'true'}}}} def test_decider(): """Test decider method.""" schedule = MagicMock() activity_result = MagicMock() activity_result.result = { 'update_dim_tables.sns_report_subject': 'Test subj', 'bootstrap.licensor': 'theorchard', } schedule.return_value = activity_result # list of activities activity_names = { 'bootstrap', 'grab_drop_files', 'update_feed_s3_file_status', 'set_feed_status_downloaded', 'clean_staging_raw', 'populate_temp_staging_tables', 'populate_staging_raw', 'mark_staging_raw_table_tasks_complete', 'set_status_to_populated_raw_table', 'drop_temp_staging_tables', 'update_dim_tables', 'send_dimension_tables_update_report_sns_notification', 'load_staging_fact_table', 'load_fact_tables', 'build_jenkins_dbt' } # test activities are called flow = Flow() flow.decider(schedule) scheduled_activities = set() for (name, activity), _ in schedule.call_args_list: assert name in activity_names assert isinstance(activity, Activity) scheduled_activities.add(name) assert scheduled_activities == activity_names # test bootstrap stops flow if data is already ingested schedule.reset_mock() activity_names = {'bootstrap'} activity_result.result = {'bootstrap.stop': True} flow.decider(schedule) scheduled_activities = set() for (name, activity), _ in schedule.call_args_list: assert name in activity_names assert isinstance(activity, Activity) scheduled_activities.add(name) assert scheduled_activities == activity_names # test update_feed_s3_file_status stops flow if files are unavailable activity_names = { 'bootstrap', 'grab_drop_files', 'update_feed_s3_file_status', } activity_result.result = { 'update_dim_tables.sns_report_subject': 'Test subj', 'bootstrap.licensor': 'theorchard', 'update_feed_s3_file_status.file_status': garcon_feed_status.STATUS_NOT_AVAILABLE } schedule.reset_mock() flow.decider(schedule) scheduled_activities = set() for (name, activity), _ in schedule.call_args_list: assert name in activity_names scheduled_activities.add(name) assert scheduled_activities == activity_names @patch('feed_ingestion.tasks.garcon_feed_status.get_status') def test_drop_files_generator(mack_get_status): """Test drop_files_generator method.""" # generated files generator = list(Flow().drop_files_generator(_CONTEXT)) # check metadata file is generated metadata_file_info = { 'source_key_name': 'data/2000/01/01/' 'orchard_metadata_2000-01-01.txt.bz2', 'destination_key_name': 'Pandora/archives/2000-01-01/' 'orchard_metadata_2000-01-01.txt.bz2'} assert metadata_file_info in generator # check country-level streams file is generated country_file_info = { 'source_key_name': 'data/2000/01/01/' 'orchard_US_2000-01-01.txt.bz2', 'destination_key_name': 'Pandora/archives/2000-01-01/' 'orchard_US_2000-01-01.txt.bz2'} assert country_file_info in generator def test_temp_staging_tables_generator(): """Test temp_staging_tables_generator method.""" # generated parameters generator = list(Flow().temp_staging_tables_generator(_CONTEXT)) # check metadata table parameters are generated metadata_params = _CONTEXT[ 'bootstrap.temp_staging_raw_tables']['metadata'] assert metadata_params in generator error_on_column_count_mismatch = _CONTEXT[ 'bootstrap.temp_staging_raw_tables'][ 'metadata']['kwargs'] for generators in generator: assert error_on_column_count_mismatch == generators['kwargs'] # check country-level streams table parameters are generated streams_params = _CONTEXT['bootstrap.temp_staging_raw_tables']['US'] assert streams_params in generator