"""CLI functions to run decider, worker and exec processes.""" import argparse from datetime import datetime from datetime import timedelta import json import logging import time import boto3 from garcon import activity from garcon import decider from yt_conflict_elasticsearch import flows logger = logging.getLogger('cli') def _is_flow_running(swf_domain, workflow_id): """Check is flow running right now. Args: swf_domain (str): SWF domain name. workflow_id (str): SWF workflow id. Returns: bool: is flow running or not. """ client = boto3.client('swf', region_name='us-east-1') response = client.count_open_workflow_executions( domain=swf_domain, startTimeFilter=dict(oldestDate=(datetime.now() - timedelta(days=10))), executionFilter=dict(workflowId=workflow_id) ) 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. """ workflow_id = flow.workflow_id(json.loads(context)) if _is_flow_running(flow.domain, workflow_id): logger.info( 'The flow "{}" is already running with workflowId "{}". Exiting.' .format(flow.name, workflow_id)) return client = boto3.client('swf', region_name='us-east-1') start_kwargs = dict( domain=flow.domain, workflowId=workflow_id, 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)) execution = client.start_workflow_execution(**start_kwargs) if json.loads(context).get('wait_until_complete') == 'True': print('Polling until Workflow is complete') max_time = datetime.now() + timedelta(seconds=flow.timeout) while datetime.now() < max_time: current_exec = client.describe_workflow_execution( domain=flow.domain, execution={ 'workflowId': workflow_id, 'runId': execution['runId'] } ) if current_exec['executionInfo'].get( 'executionStatus') != 'CLOSED': time.sleep(10) else: if current_exec['executionInfo'].get( 'closeStatus') == 'COMPLETED': return execution else: raise Exception( 'Workflow with runId {runId} in domain {domain} ' 'failed'.format( runId=execution['runId'], domain=flow.domain)) else: return execution 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 garcon(*args): """Define 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 = flows.get_flow(run_args.flow) garcon_command = _COMMANDS[run_args.cmd] garcon_command( flow=flow, context=run_args.context or '{}')