""" iTunes Garcon tasks. Tasks to ingest iTunes data. """ from collections import defaultdict from datetime import datetime from datetime import timedelta import os from botocore.exceptions import ClientError from garcon import task from garcon_contrib.dynamo_feed_status import garcon_feed_status from feed_ingestion.flows.apple_music_streams.utils import \ get_missing_vendors from feed_ingestion.flows.itunes import config from feed_ingestion.tasks import reporter_tasks from feed_ingestion.tasks.feed_status_tasks import \ get_contexts_config_for_report from feed_ingestion.tasks.s3_tasks import copy_s3_key from feed_ingestion.util import check_nonconcurrent_workflows from feed_ingestion.util import itunes_reporter from feed_ingestion.util import task_status from feed_ingestion.util.aws.s3 import get_list_of_files_and_directories from feed_ingestion.util.context_util import strtobool STOP_RESPONSE = {'stop': True} @task.decorate(timeout=900) def extract_reporter_file_to_s3( activity, reporter_account, report_type, date, destination_s3_path, feed_name=None, report_role='sales', report_country=None, vendors=None, licensor=None): """iTunes-local wrapper around common reporter extraction task. Keeps shared behavior unchanged while allowing flow-specific timeout. """ return reporter_tasks.extract_reporter_file_to_s3( activity=activity, reporter_account=reporter_account, report_type=report_type, date=date, destination_s3_path=destination_s3_path, feed_name=feed_name, report_role=report_role, report_country=report_country, vendors=vendors, licensor=licensor, ) @task.decorate(timeout=1000) def check_concurrent_status(activity, licensor, domain): """Check if there is any active executions. Args: activity (ActivityWorker): The Garcon activity worker. licensor (str): The licensor to ingest. domain (str): The flow domain (f.e. prod_feed_ingestion). Returns: dict: {} or STOP_RESPONSE from check_ingested_status decorator """ return check_nonconcurrent_workflows.check_concurrent_status( domain, licensor, config.max_concurrent_executions[licensor]) def get_version(rules, lookup_date): """Get version of the report by date.""" version, format_ = '', '.txt' for rule in rules: period_start_str = rule['since'] period_start = datetime.strptime( period_start_str, '%Y-%m-%d').date() if period_start > lookup_date: break version = rule['version'] format_ = rule['format'] return version, format_ @task.decorate(timeout=1000) def bootstrap(activity, date, reload, licensor, snowflake_error_limit=None, populate_only='False', use_s3='False', snowflake_error_on_column_count_mismatch=None ): """Bootstrap workflow by injecting initial context from config. Args: activity (ActivityWorker): The Garcon activity worker. date (str): Reporting date (YYYY-MM-DD). reload (str or None): If 'True' delete all feed statuses in DynamoDB. licensor (str): The licensor to ingest. snowflake_error_limit (int or None): Snowflake error limit. populate_only (str): If flow should finish on populating tables. use_s3 (str): If flow should backfill data from S3 only. snowflake_error_on_column_count_mismatch: Returns: dict: Initial context for the workflow. """ # date is the date passed in or yesterday's date activity.logger.info(f'Starting itunes ' f'for {config.feed_name}_{licensor}') if not date: date_obj = datetime.today() - timedelta(days=1) else: date_obj = datetime.strptime(date, '%Y-%m-%d') date_str = date_obj.strftime('%Y-%m-%d') overall_feed_name = '_'.join([config.feed_name, licensor]) if reload == 'True': activity.logger.info('Delete status for feed: {} {} '.format( overall_feed_name, date)) garcon_feed_status.delete_status(overall_feed_name, date) else: overall_status = garcon_feed_status.get_overall_status( overall_feed_name, date) completed = task_status.is_completed_overall_job( overall_feed_name, date) if overall_status == garcon_feed_status.STATUS_INGESTED and completed: activity.logger.info('Job is already completed.') return STOP_RESPONSE vendors_config = get_vendors_config(licensor) populate_only = populate_only if populate_only is not None else 'False' use_s3 = use_s3 if use_s3 is not None else 'False' reporter_accounts = {} for reporter_account in config.licensors[licensor]: reporter = itunes_reporter.get_reporter( reporter_account, date, vendors=vendors_config) filename = reporter.get_sales_report_file_name(config.report_name) if licensor == 'awal' and use_s3 == 'True': version, format_ = get_version( config.awal_s3_versions['versions'], date_obj.date()) filename = config.awal_file_pattern.format( vendor=config.licensors['awal'][0], date=date_obj, version=version, format=format_) # useless now reporter_accounts[reporter_account] = filename garcon_feed_status.set_missing_files(overall_feed_name, date, []) task_status.delete_newcontexts(overall_feed_name, date) # init contexts, if they don't exist contexts_field_value = task_status.get_report_contexts( overall_feed_name, date) if not contexts_field_value: activity.logger.info(f'Creating contexts for {overall_feed_name}') conf_optional, conf_contexts = ( get_contexts_config_for_report( 'default', config.contexts_config, licensor)) vendors = list(config.licensors[licensor]) task_status.create_report_contexts( overall_feed_name, date, vendors, conf_contexts, conf_optional) if isinstance(snowflake_error_limit, int): snowflake_error_limit = snowflake_error_limit else: snowflake_error_limit = config.snowflake_error_limit snowflake_eoccm = snowflake_error_on_column_count_mismatch or 'True' populate_only = bool(strtobool(populate_only)) return { 'date': date_str, 'feed_name': overall_feed_name, 'licensor': licensor, 'processed_datetime': datetime.now().isoformat(), 'reporter_accounts': reporter_accounts, 's3_archive_bucket': config.s3['archive_bucket'].format( date=date), 'snowflake_error_limit': snowflake_error_limit, 'expected_download_files': list(reporter_accounts.values()), 'populate_only': populate_only, 'use_s3': bool(strtobool(use_s3)), 'snowflake_error_on_column_count_mismatch': snowflake_eoccm, 'contexts': config.contexts_config } def get_vendors_config(licensor): """Map vendor ids and property files for all vendors. Args: licensor (str): The licensor to ingest. Returns: dict: with vendor ids and names of property files. """ vendors = defaultdict(dict) for vendor_id in config.licensors[licensor]: vendor_account = config.vendor_account_mapping.get( vendor_id, config.vendor_account_mapping.get(licensor)) vendors[vendor_id] = { 'VENDOR_ID': vendor_id, 'licensor': licensor, 'ACCOUNT': vendor_account } return vendors @task.decorate(timeout=36000) def grab_drop_files_from_s3_awal( activity, feed_name, date): """Upload files to archive location. Args: activity (ActivityWorker): The activity worker. feed_name (str): Feed name of workflow execution for status updates. date (str): Reporting date (YYYY-MM-DD). reporter_account (str): Apple Reporter Account ('ORCHARD' or 'IODA'). """ parsed_date = datetime.strptime(date, '%Y-%m-%d').date() from_path = config.awal_s3_path.format( drop_bucket=config.awal_drop_bucket, s3_folder=config.awal_s3_versions['s3_folder'], version=config.awal_version, date=parsed_date) s3_archive_path = config.s3['archive_bucket'].format(date=date) version, format_ = get_version( config.awal_s3_versions['versions'], parsed_date) filename = config.awal_file_pattern.format( vendor=config.licensors['awal'][0], date=parsed_date, version=version, format=format_) files = get_list_of_files_and_directories(from_path) if not files: return {filename: garcon_feed_status.STATUS_NOT_AVAILABLE} file_ = '' for file in files: if filename in file: file_ = file # there is usually only one file file = os.path.basename(file_) to_path = f'{s3_archive_path}{file}' try: copy_s3_key(f'{from_path}{file}', to_path) activity.logger.info( 'File was copied from {from_path} to {to_path}'.format( from_path=from_path, to_path=to_path)) except ClientError as err: activity.logger.error( ('Cannot copy files from {drop_location} ' 'to {archive_location}. {exception_body}').format( drop_location=from_path, archive_location=to_path, exception_body=err)) return {filename: garcon_feed_status.STATUS_NOT_AVAILABLE} garcon_feed_status.set_status( feed_name, date, filename, status=garcon_feed_status.STATUS_DOWNLOADED) # add current vendor to set of new_contexts vendor_id = config.licensors['awal'][0] task_status.add_newcontext(feed_name, date, vendor_id) task_status.update_report_context_status( feed_name, date, vendor_id, task_status.CONTEXT_STATUS_IN_PROGRESS) return {filename: garcon_feed_status.STATUS_DOWNLOADED} @task.decorate(timeout=600) def check_available_reports(activity, date, s3_path, licensor): """Return new vendors which are available. Args: activity (ActivityWorker): The Garcon activity worker. date (str): Reporting date (YYYY-MM-DD). licensor (str): The licensor to ingest. s3_path: is required in order to better parse missing_files field Returns: dict: of available vendors. if report is missing for optional vendor than file will be ignored. """ feed_name = f'{config.feed_name}_{licensor}' in_progress_contexts = task_status.get_report_in_progress_contexts( feed_name, date) if (task_status.is_completed_task(feed_name, date, 'reporter_to_s3') and in_progress_contexts): missing_vendors = get_missing_vendors(feed_name, date, s3_path) available_reports = [vendor for vendor in config.licensors[licensor] if vendor not in missing_vendors] activity.logger.info(f'Available vendors: {available_reports}') return {'available_reports': available_reports} else: activity.logger.info('Skip processing, ' 'because there are no new vendors') return STOP_RESPONSE