"""CLI functions to run decider, worker and exec processes.""" import argparse from datetime import datetime from datetime import timedelta import json import importlib import logging import time import boto3 from garcon import activity from garcon import decider from garcon_boilerplate import flows logger = logging.getLogger('cli') def _flow_is_running(flow): """Check is flow running right now. Args: flow (object): Garcon flow object. Returns: bool: is flow running or not. """ client = boto3.client('swf', region_name='us-east-1') response = client.count_open_workflow_executions( domain=flow.domain, startTimeFilter={ 'oldestDate': (datetime.now() - timedelta(days=10)), }, typeFilter={ 'name': flow.name, 'version': flow.version }) return response['count'] > 0 def execute_flow(flow, context, **kwargs): """Launch the workflow execution. Args: flow (object): Garcon flow module. context (str): Initial context parsed from JSON. output_file (str): Optional path to file to store all args as well as workflow_id and run_id. kwargs: Extensible call API. """ if _flow_is_running(flow) : logger.info('The flow "{}" is running.'.format(flow.name)) return client = boto3.client('swf', region_name='us-east-1') start_kwargs = dict( domain=flow.domain, workflowId=flow.workflow_id(json.loads(context)), workflowType=dict( name=flow.name, version=flow.version), taskList=dict(name=flow.name), input=context) if hasattr(flow, 'timeout'): start_kwargs.update(executionStartToCloseTimeout=str(flow.timeout)) return client.start_workflow_execution(**start_kwargs) def run_decider(flow, **kwargs): """Launch the SWF decider process. Args: flow (module): Garcon flow module. """ garcon_decider = decider.DeciderWorker(flow) while True: garcon_decider.run() time.sleep(1) def run_activity_worker(flow, **kwargs): """Launch the activity worker process. Args: flow (module): Garcon flow module. """ garcon_worker = activity.ActivityWorker(flow) garcon_worker.run() _COMMANDS = { 'exec': execute_flow, 'decider': run_decider, 'worker': run_activity_worker} def _parse_run_args(cli_args): """Parse command line run params. Args: tuple: Args for the flow. Returns: argparse.Namespace: Object with parsed arguments. """ 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']) return parser.parse_args(cli_args) if cli_args else parser.parse_args() def _get_flow_class(flow_name): """Retrieve flow class. Args: flow_name (str): Name of a garcon flow in the project. Returns: Flow: Flow class if found. Raises: AssertionError: If the flow implementation wasn't found. """ # import the flow module flow_module = importlib.import_module( '.{}.flow'.format(flow_name), flows.__name__) # look for the flow class flow_class = getattr(flow_module, 'Flow', None) assert flow_class, 'Flow "{}" implementation was not found.'.format( flow_name) return flow_class def garcon(*args): """Main entry point for the Garcon command line integration. Args: args (tuple): Args for the flow. """ run_args = _parse_run_args(args) logging.basicConfig(level=getattr(logging, run_args.log_level.upper())) flow_class = _get_flow_class(run_args.flow) flow = flow_class() garcon_command = _COMMANDS[run_args.cmd] garcon_command( flow=flow, context=run_args.context or '{}')