import argparse import importlib import json import logging import time import boto.swf.layer2 as swf from garcon import activity from garcon import decider from processing_accounting import flows from processing_accounting.util import cli as cli_util def execute_flow(flow, context, **kwargs): """Launches the workflow execution. Args: flow (module): garcon flow module context (str): initial context parsed from json kwargs: extensible call api """ start_kwargs = dict( workflow_id=flow.workflow_id(json.loads(context)), input=context) if hasattr(flow, 'timeout'): start_kwargs.update(execution_start_to_close_timeout=str(flow.timeout)) if cli_util.check_required_params(flow, context) is False: print( 'Missing required param(s): {}'.format( ', '.join(flow.required_params))) print( 'Available params are: {}'.format( ', '.join(flow.available_params))) return return swf.WorkflowType( name=flow.name, domain=flow.domain, version=flow.version, task_list=flow.name).start(**start_kwargs) def run_decider(flow, **kwargs): """Launches the SWF decider process. Args: flow (module): garcon flow module kwargs: extensible call api """ worker = decider.DeciderWorker(flow) while True: worker.run() time.sleep(1) def run_activity_worker(flow, **kwargs): """Launches the activity worker process. Args: flow (module): garcon flow module kwargs: extensible call api """ worker = activity.ActivityWorker(flow) worker.run() _COMMANDS = {'exec': execute_flow, 'decider': run_decider, 'worker': run_activity_worker} def garcon(*args): """Main entry point for the Garcon command line integration. """ parser = argparse.ArgumentParser(description='Garcon command line util') parser.add_argument('cmd', choices=_COMMANDS.keys(), help='garcon command') parser.add_argument('flow', choices=flows.__all__, help='name of the flow') parser.add_argument('-c', '--context', help='initial context [json]') parser.add_argument('-l', '--log_level', help='logging level', default='error', choices=['critical', 'error', 'warning', 'info', 'debug']) # parse cl args (allows easy unit testing) args = parser.parse_args(args) if args else parser.parse_args() # import the flow module flow = importlib.import_module( '.{}.flow'.format(args.flow), flows.__name__) # look for the flow class, if it exists, otherwise fall back # on the flow module itself. flow_class = getattr(flow, 'Flow', None) # execute command args.context = args.context or '{}' logging.basicConfig(level=getattr(logging, args.log_level.upper())) if args.cmd == 'exec': execute_flow(flow=flow_class(), context=args.context) elif args.cmd == 'decider': run_decider(flow=flow_class(), context=args.context) elif args.cmd == 'worker': run_activity_worker(flow=flow_class(), context=args.context) else: raise Exception('Option is not available.')