"""CLI functions to run decider, worker and exec processes.""" from argparse import ArgumentParser from datetime import date import json import sys import boto.swf.layer2 as swf from garcon import activity from garcon import decider from job.flow import SnowflakeUnload flow = SnowflakeUnload() parser = ArgumentParser( description='Analytics Worker command line.', epilog='© {year} The Orchard'.format(year=str(date.today().year))) parser.add_argument( '-s', '--start', dest='start', help='start a flow.', action='store_true') parser.add_argument( '-c', '--context', dest='context', help='flow context', default='{}') parser.add_argument( '-w', '--worker', dest='worker', help='define the type of worker (activity, decider).') parser.add_argument( '-a', '--activities', dest='activities', default=None, help='define the name of the activity (activity, decider).') def start(context=None): """Start a flow. Args: context (dict): the context to pass to the flow. Returns: WorkflowExecution: executing the workflow returns information about the execution. """ workflow_id = flow.workflow_id(context) return swf.WorkflowType( name=flow.name, domain=flow.domain, version='3.0', task_list=flow.name).start( workflow_id=workflow_id, execution_start_to_close_timeout='10800', # 3 hours. input=json.dumps(context or {})) if __name__ == '__main__': options = parser.parse_args() if len(sys.argv) <= 1: parser.print_help() exit() if options.start: start(json.loads(options.context)) else: if options.worker == 'decider': w = decider.DeciderWorker(flow) while(True): w.run() elif options.worker == 'activity': w = activity.ActivityWorker(flow, options.activities) w.run()