"""FTP-related tasks.""" from garcon import task from garcon_contrib.dynamo_feed_status import garcon_feed_status from garcon_contrib.ftp import garcon_ftp from feed_ingestion.flows.mrc import config @task.decorate(timeout=28800) def fetch_from_drop_location( activity, date, source_files_dict, s3_archive_path, ftp_creds, feed_name=None, need_all_files=True, mapping=False): """Copy source file from FTP to S3. Args: activity (Activity): Activity instance. date (str): date being processed. source_files_dict (dict): dict containing metadata of files being processed: { 'files": [ { 'file_name': 'my_cute_file' 'required': True }, { 'file_name': 'my_nasty_file' }, ] }. s3_archive_path (str): S3 directory to write files to. ftp_creds (dict): FTP credentials. feed_name (str): Name of the feed. need_all_files (bool): If True, all the files need to be present. Returns: dict: Context patch. """ def copy_file(file_name, s3_archive_path, date=None): """Copy one file from FTP to S3.""" if file_name[:file_name.rfind('_')] in \ config.mapping_snowflake_table_suffixes: s3_archive_path = s3_archive_path.format( mapping=file_name[:file_name.rfind('_')], date=date) copy_response = garcon_ftp.copy_file_from_ftp_to_s3( activity, ftp_creds, ftp_creds['path'], file_name, s3_archive_path, file_name) found = copy_response['status'] if found: activity.logger.info('File downloaded: %s', copy_response) else: activity.logger.info('File not found: %s', copy_response) return { 'file_name': file_name, 'file_size': copy_response.get('file_size', -1), 'found': found } required_files = set() if source_files_dict: activity.logger.info('Fetching source files: %s', date) required_files = { fd['file_name'] for fd in source_files_dict['files'] if fd.get('required', False) or need_all_files } result = [ copy_file(fd['file_name'], s3_archive_path) for fd in source_files_dict['files'] ] elif mapping: activity.logger.info('Fetching mapping files: %s', date) source_files_dict = garcon_ftp._get_list_of_files_and_directories_sftp( ftp_creds, ftp_creds['path']) result = [copy_file(fd, s3_archive_path, date) for fd in source_files_dict] missing_files = { fd['file_name'] for fd in result if not fd['found'] } if required_files & missing_files: if feed_name is not None: garcon_feed_status.set_overall_status( feed_name, date, garcon_feed_status.STATUS_NOT_INGESTED) activity.logger.error('Failed to download files from FTP') return dict(stop=True) return dict(source_files_dict={'files': result})