"""Unit tests for QQ 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.qq.flow import Flow _CONTEXT = { 'bootstrap.date': '2000-01-01', 'bootstrap.drop_path': 's3://drop_bucket/data/', 'bootstrap.s3_full_path': 's3://data_bucket/qq/archives/2000-01-01/', 'bootstrap.expected_files': ['P926_20000101_20000101_TheOrchard.txt']} def test_decider(): """Test decider method.""" schedule = MagicMock() activity_result = MagicMock() activity_result.result = { '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', 'source_files', 'load_staging_raw', 'mark_staging_raw_table_tasks_complete', 'set_status_to_populated_raw_table', 'load_staging_fact_table', 'collect_kwargs_activity', 'load_fact_table' } # 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 = { '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(mock_get_status): """Test drop_files_generator method.""" # generated files generator = list(Flow().drop_files_generator(_CONTEXT)) # check whether file is generated data_file_info = { 'source_key_name': 'data/P926_20000101_20000101_TheOrchard.txt', 'destination_key_name': 'qq/archives/2000-01-01/P926_20000101_20000101_TheOrchard.txt'} assert data_file_info in generator