"""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 sfn_arn = f'arn:aws:states:us-east-1:437795906767:stateMachine:{env}-sr-fingerprinter-sfn' # noqa:E501 sfn_client = boto3.client('stepfunctions') # read in CSV file and validate assets = [] with open(filename, 'r') as f: reader = csv.DictReader(f) if set(reader.fieldnames) != set(['ASSET_ID', 'FILENAME', 'EXTENSION']): # noqa:E501 print('Invalid input CSV column format!') exit(1) assets = [x for x in reader] # sort to allow re-processing from middle process = False assets = sorted(assets, key=lambda x: x['ASSET_ID']) for asset in assets: asset_id = asset['ASSET_ID'] # skip if processing only after specific assset if not process and after: if asset_id == after: process = True 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( [ 'OrchardAsset', asset_id, 'created', str(int(time.time())) ] ) payload = { 'id': asset_id, 'label': 'OrchardAsset', 'operation': 'created', 'payload': { 'after': { 'properties': { 'id': asset_id, 'filename': asset['FILENAME'], 'extension': asset['EXTENSION'] } } } } print(exc_name) sfn_client.start_execution( stateMachineArn=sfn_arn, name=exc_name, input=json.dumps(payload) ) if __name__ == '__main__': main()