"""Docstring.""" import csv import json import sys import time import boto3 def main(): """Entrypoint.""" filename = sys.argv[1] env = sys.argv[2] after = sys.argv[3] if len(sys.argv) > 3 else None max_running = 15 print('filename', filename) print('env', env) print('after', after) sfn_arn = f'arn:aws:states:us-east-1:437795906767:stateMachine:{env}-sr-fingerprinter-sfn' # noqa:E501 sfn_client = boto3.client('stepfunctions', region_name='us-east-1') # read in CSV file and validate osrs = [] with open(filename, 'r') as f: reader = csv.DictReader(f) if set(reader.fieldnames) != set(['OSR_ID']): # noqa:E501 print('Invalid input CSV column format!') exit(1) osrs = [x for x in reader] # sort to allow re-processing from middle process = False osrs = sorted(osrs, key=lambda x: x['OSR_ID']) num_processed = 0 total_osrs = len(osrs) for osr in osrs: osr_id = osr['OSR_ID'] # skip if processing only after specific osr if not process and after: if osr_id == after: process = True num_processed += 1 continue # throttle based on num running state machines while True: sfn_list_response = sfn_client.list_executions( stateMachineArn=sfn_arn, statusFilter='RUNNING', maxResults=max_running ) num_running = len(sfn_list_response['executions']) if num_running >= max_running: time.sleep(5) else: break # kick-off new execution exc_name = '-'.join( [ 'OrchardSoundRecording', osr_id, 'updated', str(int(time.time())) ] ) payload = { 'id': osr_id, 'label': 'OrchardSoundRecording', 'operation': 'updated', 'payload': {} } print(exc_name) sfn_client.start_execution( stateMachineArn=sfn_arn, name=exc_name, input=json.dumps(payload) ) num_processed += 1 percent_processed = (num_processed / total_osrs) * 100 print(f'Processed {osr_id} {num_processed} / {total_osrs} : {percent_processed}%') # noqa:E501 if __name__ == '__main__': main()