"""Retry past execution.""" import argparse import json import time import boto3 import botocore MAX_RUNNING = 1000 def main(): """Entrypoint.""" args = parse_args() filename = args.filename suffix = args.suffix sleep = args.sleep wake = args.wake after = args.after # validate throttle args if sleep <= wake: raise Exception('sleep <= wake') if sleep > MAX_RUNNING: raise Exception('sleep > MAX') if sleep < 1: raise Exception('sleep < 1') if wake < 0: raise Exception('wake < 0') sfn_client = boto3.client('stepfunctions') # expected format from ../sfn-execution-searcher/search.py with open(filename, 'r') as f: data = f.read() logs = json.loads(data) sfn_arn = logs['arn'] executions = logs['executions'] # setup skip over some records skip = after is not None # retry each execution, adding suffix to name for execution in executions: if execution['name'].endswith('retry1'): continue # skip processing until after param is seen if skip and execution['name'] == after: skip = False continue if skip: continue # query for details about execution to retry details = sfn_client.describe_execution( executionArn=execution['arn'] ) exc_name = details['name'] retry_name = details['name'] + suffix try: sfn_client.start_execution( stateMachineArn=sfn_arn, name=retry_name, input=details['input'] ) print(f'Started: {retry_name}') # determine if script should slow down running_count = _num_running(sfn_client, sfn_arn) if running_count >= sleep: while True: print(f'Cooling down from {running_count} concurrency') running_count = _num_running(sfn_client, sfn_arn) if running_count < wake: print(f'Warming up from {running_count} concurrency') # noqa:E501 break time.sleep(5) except botocore.exceptions.ClientError as e: if e.response['Error']['Code'] == 'ExecutionAlreadyExists': print(f'Already retried: {exc_name}') elif e.response['Error']['Code'] == 'InvalidName': print(f'Invalid name: {retry_name}') else: raise e print('Done!') def _num_running(sfn_client, sfn_arn): sfn_list_response = sfn_client.list_executions( stateMachineArn=sfn_arn, statusFilter='RUNNING', maxResults=MAX_RUNNING ) return len(sfn_list_response['executions']) def parse_args(): """Read CLI args.""" parser = argparse.ArgumentParser( formatter_class=argparse.ArgumentDefaultsHelpFormatter) parser.add_argument( 'filename', help='SFN executions from ../sfn-execution-searcher' ) parser.add_argument( 'suffix', help='String to append to end of SFN execution name on retry' ) parser.add_argument( '--sleep', type=int, default=MAX_RUNNING, help='Pause new executions if num running >=' ) parser.add_argument( '--wake', type=int, default=1, help='Start new executions if num running <' ) parser.add_argument( '--after', type=str, help='Only retry after SFN execution name' ) return parser.parse_args() if __name__ == '__main__': main()