"""PROPER Outgoing Feed for Changed Releases. This generates the Changed Releases feed for PROPER. """ from garcon import runner from garcon.param import StaticParam from feed_sender.flows import helpers from feed_sender.flows import tasks as base_tasks from feed_sender.flows.base import FlowBase from feed_sender.flows.proper_changed_releases import config from feed_sender.flows.proper_changed_releases import tasks from feed_sender.garcon_contrib.aws import garcon_s3 as s3utils from feed_sender.garcon_contrib.ftp import garcon_ftp class Flow(FlowBase): """Class representing the workflow.""" def __init__(self): """Initialise flow object.""" super(Flow, self).__init__( feed_name=config.FEED_NAME, version='1.0') self.timeout = 1800 def decider(self, schedule, context=None): """Activity decider. Args: schedule (callable): The scheduler method. context (dict): Initial context of the workflow. """ bootstrap = schedule('bootstrap', self.bootstrap) # Stop flow if feed already successfully ran and reload is False. if bootstrap.result.get('bootstrap.stop'): return if context.get('proper_in_vector'): clean_up_old_files = schedule( 'clean_up_old_files', self.clean_up_old_files, requires=[bootstrap]) set_status_to_processing = schedule( 'set_status_to_processing', self.set_status_to_processing, requires=[clean_up_old_files]) obtain_jobs = schedule( 'obtain_jobs', self.obtain_jobs, requires=[set_status_to_processing]) # if there are no jobs, complete swf job and return message if obtain_jobs.result.get('obtain_jobs.error'): schedule( 'set_status_to_completed', self.set_status_to_completed) return obtain_jobs.result.get('obtain_jobs.error') set_job_status_to_encoding = schedule( 'set_job_status_to_encoding', self.set_job_status_to_encoding, requires=[obtain_jobs]) generate_feed_for_changed_releases = schedule( 'generate_feed_for_changed_releases', self.generate_feed_for_changed_releases, requires=[set_job_status_to_encoding]) set_job_status_to_encoded = schedule( 'set_job_status_to_encoded', self.set_job_status_to_encoded, requires=[generate_feed_for_changed_releases]) set_job_status_to_delivering = schedule( 'set_job_status_to_delivering', self.set_job_status_to_delivering, requires=[set_job_status_to_encoded]) # if s3_only flag passed in context, skip ftp/sns activities if 's3_only' in context: schedule( 'set_status_to_uploaded_to_s3', self.set_status_to_uploaded_to_s3, requires=[set_job_status_to_delivering]) return transport_to_sftp = schedule( 'transport_to_sftp', self.transport_to_sftp, requires=[generate_feed_for_changed_releases]) if transport_to_sftp.result.get('upload_to_sftp.error'): schedule( 'set_job_status_to_delivery_failure', self.set_job_status_to_delivery_failure, requires=[transport_to_sftp]) else: schedule( 'set_job_status_to_delivered', self.set_job_status_to_delivered, requires=[transport_to_sftp]) # # Update status on Dynamo DB schedule( 'set_status_to_completed', self.set_status_to_completed, requires=[transport_to_sftp]) else: clean_up_old_files = schedule( 'clean_up_old_files', self.clean_up_old_files, requires=[bootstrap]) set_status_to_processing = schedule( 'set_status_to_processing', self.set_status_to_processing, requires=[clean_up_old_files]) generate_feed_for_changed_releases = schedule( 'generate_feed_for_changed_releases', self.generate_feed_for_changed_releases, requires=[set_status_to_processing]) # if s3_only flag passed in context, skip ftp/sns activities if 's3_only' in context: schedule( 'set_status_to_uploaded_to_s3', self.set_status_to_uploaded_to_s3, requires=[generate_feed_for_changed_releases]) return transport_to_ftp = schedule( 'transport_to_ftp', self.transport_to_ftp, requires=[generate_feed_for_changed_releases]) update_delivery_history = schedule( 'update_delivery_history', self.update_delivery_history, requires=[transport_to_ftp]) # # Update status on Dynamo DB schedule( 'set_status_to_completed', self.set_status_to_completed, requires=[update_delivery_history]) @property def bootstrap(self): """Bootstrap initial configuration.""" return self.create( name='bootstrap', tasks=runner.Sync( tasks.bootstrap.fill( namespace='bootstrap', cutoff='cutoff', reload='reload', s3_only='s3_only'))) @property def clean_up_old_files(self): """Clean up old new releases file.""" return self.create( name='clean_up_old_files', retry=5, tasks=runner.Sync( s3utils.remove_files_from_path.fill( namespace='clean_up_old_files_on_s3', path='bootstrap.changed_releases_s3_path'))) @property def set_job_status_to_encoding(self): """Set job status to Encoding.""" return self.create( name='set_job_status_to_encoding', tasks=runner.Sync( tasks.update_job_status.fill( namespace='set_job_status_to_encoding', s3path='bootstrap.changed_releases_s3_path', job_ids_filename='bootstrap.job_ids_filename', status=StaticParam('encoding')))) @property def set_job_status_to_encoded(self): """Set job status to Encoding.""" return self.create( name='set_job_status_to_encoded', tasks=runner.Sync( tasks.update_job_status.fill( namespace='set_job_status_to_encoded', s3path='bootstrap.changed_releases_s3_path', job_ids_filename='bootstrap.job_ids_filename', status=StaticParam('encoded')))) @property def set_job_status_to_delivering(self): """Set job status to Encoding.""" return self.create( name='set_job_status_to_delivering', tasks=runner.Sync( tasks.update_job_status.fill( namespace='set_job_status_to_delivering', s3path='bootstrap.changed_releases_s3_path', job_ids_filename='bootstrap.job_ids_filename', status=StaticParam('delivering')))) @property def set_job_status_to_delivered(self): """Set job status to Encoding.""" return self.create( name='set_job_status_to_delivered', tasks=runner.Sync( tasks.update_job_status.fill( namespace='set_job_status_to_delivered', s3path='bootstrap.changed_releases_s3_path', job_ids_filename='bootstrap.job_ids_filename', status=StaticParam('delivered')))) @property def set_job_status_to_delivery_failure(self): """Set job status to Encoding.""" return self.create( name='set_job_status_to_delivery_failure', tasks=runner.Sync( tasks.update_job_status.fill( namespace='set_job_status_to_delivery_failure', s3path='bootstrap.changed_releases_s3_path', job_ids_filename='bootstrap.job_ids_filename', status=StaticParam('delivery_failure'), error_log='upload_to_sftp.error'))) @property def generate_feed_for_changed_releases(self): """Get PROPER Changed Releases.""" return self.create( name='generate_feed_for_changed_releases', tasks=runner.Sync( tasks.generate_feed_for_changed_releases.fill( namespace='generate_feed_for_changed_releases', cutoff='bootstrap.cutoff', s3path='bootstrap.changed_releases_s3_path', filename='bootstrap.changed_releases_filename', release_ids_filename='bootstrap.release_ids_filename', proper_in_vector='proper_in_vector', job_ids_filename='bootstrap.job_ids_filename'))) @property def transport_to_ftp(self): """Transport files from S3 to FTP.""" return self.create( name='transport_to_ftp', retry=5, tasks=runner.Sync( garcon_ftp.copy_from_s3_to_ftp.fill( sftp_creds=StaticParam(config.SFTP_CREDS), file_names_list=( 'bootstrap.filenames_for_sftp_activity'), remote_s3_dir_path=( 'bootstrap.changed_releases_s3_path'), remote_ftp_dir_path=( 'bootstrap.changed_releases_sftp_path')))) @property def transport_to_sftp(self): """Transport files from S3 to FTP and set status UPLOADED_TO_FTP.""" return self.create( name='transport_to_sftp', tasks=runner.Sync( tasks.upload_to_sftp.fill( namespace='upload_to_sftp', sftp_creds=StaticParam(config.SFTP_CREDS), file_names_list=( 'bootstrap.filenames_for_sftp_activity'), remote_s3_dir_path=( 'bootstrap.changed_releases_s3_path'), remote_ftp_dir_path=( 'bootstrap.changed_releases_sftp_path')))) @property def update_delivery_history(self): """Update History table and Status in Dynamo DB.""" return self.create( name='update_delivery_history', tasks=runner.Sync( tasks.update_delivery_history.fill( namespace='update_delivery_history', cutoff='bootstrap.cutoff', s3path='bootstrap.changed_releases_s3_path', release_ids_filename='bootstrap.release_ids_filename'))) @property def set_status_to_processing(self): """Set status to PROCESSING.""" return self.create( name='set_status_to_processing', tasks=runner.Sync( base_tasks.set_status.fill( namespace='set_status_to_processing', feed_name='bootstrap.feed_name', context_date='bootstrap.context_date', status=StaticParam(helpers.STATUS_PROCESSING)))) @property def set_status_to_uploaded_to_s3(self): """Update Status in Dynamo DB when files are uploaded to S3.""" return self.create( name='set_status_to_uploaded_to_s3', tasks=runner.Sync( base_tasks.set_status.fill( namespace='set_status_to_uploaded_to_s3', feed_name='bootstrap.feed_name', context_date='bootstrap.context_date', status=StaticParam(helpers.STATUS_UPLOADED_TO_S3)))) @property def set_status_to_completed(self): """Update Status in Dynamo DB to completed.""" return self.create( name='set_status_to_completed', tasks=runner.Sync( base_tasks.set_status.fill( namespace='set_status_to_completed', feed_name='bootstrap.feed_name', context_date='bootstrap.context_date', status=StaticParam(helpers.STATUS_SENT)))) @property def obtain_jobs(self): """Obtain jobs from SQS queue.""" return self.create( name='obtain_jobs', tasks=runner.Sync( tasks.obtain_jobs.fill( namespace='obtain_jobs', s3path='bootstrap.changed_releases_s3_path', job_ids_filename='bootstrap.job_ids_filename')))