"""Generator functions for the iTunes workflow.""" from feed_ingestion.flows.itunes import config, tasks def itunes_reports_generator(context): """Generate params that we want to process. Used by reporter_to_s3 activity. Args: context (dict): The current context. Must have a 'bootstrap.reporter_accounts'. Yields: (dict): Dictionary for each reporter account / report we want to download and process. """ vendors_config = tasks.get_vendors_config(context['licensor']) for reporter_account in context['bootstrap.reporter_accounts']: yield dict( reporter_account=reporter_account, vendors_config=vendors_config) def report_data_generator(context): """Generate params for every report in reporter_accounts. Args: context (dict): The current context. Must have a 'bootstrap.reporter_accounts'. Yields: (dict): Dictionary for each reporter account / report we want to download and process. """ for vendor_id, filename in context['bootstrap.reporter_accounts'].items(): report_data = dict( vendor_id=vendor_id, temp_table=config.temp_table_name.format( licensor=context['licensor'], reporter_account=vendor_id, date=context['bootstrap.date'].replace('-', '') ), destination_s3_path=''.join(( context['bootstrap.s3_archive_bucket'], filename )), filename=filename, ) yield {'report_data': report_data} def available_vendors_generator(context): """Generate only available vendors. Args: context (dict): the current context. Must have a 'check_available_reports.available_reports' key. Yields: (dict): Dictionary for each vendor with corresponding tables info. """ available_vendors = context['check_available_reports.available_reports'] reports = report_data_generator(context) for report in reports: if report['report_data']['vendor_id'] in available_vendors: yield report