"""Snapchat Data Ingestion Workflow.""" from datetime import datetime from garcon import task from garcon_contrib.dynamo_feed_status import garcon_feed_status from snowflake_connector.etl_connector import SQLLoader from feed_ingestion.flows.snapchat import config from feed_ingestion.tasks import STOP_RESPONSE sql_loader = SQLLoader(__file__) @task.decorate(timeout=1000) def bootstrap(activity, date, report, reload): """Bootstrap workflow by getting the correct configurations. Args: activity (ActivityWorker): The activity worker. date (str): Reporting date (YYYY-MM-DD). Returns: dict: Context. """ date_obj = (datetime.strptime( date, '%Y-%m-%d') if date else datetime.today()) activity.logger.info('Bootstrap flow: {}'.format(date_obj)) report_feed_name = config.feed_name + '_' + report if reload == 'True': activity.logger.info( 'Delete status for feed: {} {} '.format(report_feed_name, date)) garcon_feed_status.delete_status(report_feed_name, date) else: overall_status = garcon_feed_status.get_overall_status( report_feed_name, date) if overall_status == garcon_feed_status.STATUS_INGESTED: return STOP_RESPONSE drop_path = config.s3.get('drop').get('path').format(date=date_obj) archive_path = config.s3.get('archive').get('path').format(date=date_obj) file_name = config.reports.get(report).get('file_name').format( date=date_obj) source_key_name = '{drop_path}/{file_name}'.format(drop_path=drop_path, file_name=file_name) destination_key_name = '{archive_path}/{file_name}'.format( archive_path=archive_path, file_name=file_name) replace_archive_files = True s3_full_path = 's3://{data_bucket}/'.format( data_bucket=config.s3.get('archive').get( 'bucket')) + archive_path + '/' return dict( date=date_obj.strftime('%Y-%m-%d'), archive_path=archive_path, source_key_name=source_key_name, destination_key_name=destination_key_name, report_feed_name=report_feed_name, replace_archive_files=replace_archive_files, s3_full_path=s3_full_path, expected_files=[file_name], staging_raw_table=config.reports.get(report).get('staging_raw_table'), file_pattern=config.reports.get(report).get('file_pattern') )