import json import os from sys import exit as sysexit import pandas import sfn # Map for event fields - {task type: detail field name} task_detail_map = { 'ExecutionFailed': 'executionFailedEventDetails', 'MapStateEntered': 'stateEnteredEventDetails', 'ChoiceStateEntered': 'stateEnteredEventDetails', 'ChoiceStateExited': 'stateExitedEventDetails', 'MapIterationSucceeded': 'mapIterationSucceededEventDetails', 'TaskStateExited': 'stateExitedEventDetails', 'LambdaFunctionStarted': None, 'MapIterationStarted': 'mapIterationStartedEventDetails', 'LambdaFunctionScheduled': 'lambdaFunctionScheduledEventDetails', 'ExecutionStarted': 'executionStartedEventDetails', 'LambdaFunctionFailed': 'lambdaFunctionFailedEventDetails', 'LambdaFunctionSucceeded': 'lambdaFunctionSucceededEventDetails', 'MapStateStarted': 'mapStateStartedEventDetails', 'TaskStateEntered': 'stateEnteredEventDetails' } SFN_NAME = os.environ.get('SFN_NAME') SFN_ARN = os.environ.get('SFN_ARN') SFN_EXEC_NAME = os.environ.get('SFN_EXEC_NAME') SFN_EXEC_ARN = \ SFN_ARN.replace('stateMachine', 'execution') + ':' + SFN_EXEC_NAME if not SFN_ARN: sysexit('You must export SFN_ARN in your environment with the SFN ARN.') if not SFN_NAME: sysexit('You must export SFN_NAME in your environment with the SFN name.') if not SFN_EXEC_NAME: sysexit('You must export SFN_EXEC_NAME in your environment with the SFN ' 'execution name.') def process_exec_list(exec_list): """Prune a list of executions to the important information.""" # Smooth out datatypes exec_dict = json.loads( json.dumps(exec_list, indent=4, sort_keys=True, default=str)) # Creat csv csv = list() print('Processing rows.') for event in exec_dict: row = get_row_from_event(event) csv.append(row) return csv def get_row_from_event(event): """Generate a CSV-shaped row from a single SFN event.""" detail_key = task_detail_map.get(event['type']) error_detail = None cause_detail = None input_row = dict() # Default no payload data name = None # Process input/output payload data if the event has a payload if detail_key: details = event[detail_key] # Load vars input_detail = json.loads(details.get('input', '{}')) output_detail = json.loads(details.get('output', '{}')) cause_detail = details.get('cause') error_detail = details.get('error') name = details.get('name') if input_detail: input_row = parse_input_field(input_detail) elif output_detail: input_row = parse_input_field(output_detail) main_row = { 'exec_order_id': event['id'], 'prev_exec_order_id': event['previousEventId'], 'timestamp': event['timestamp'], 'task_name': name, 'type': event['type'], 'error': error_detail, 'cause': cause_detail } # Combine fields row = { **main_row, **input_row } return row # Override this function to process different payload shapes def parse_input_field(detail): """Parse input/output field and return a formatted dict of values.""" input_block = dict(asset_field=None) # Assign vals input_block['product_detail'] = detail.get('product', {}) input_block['correlation_id'] = detail.get('correlation_id') input_block['asset_type'] = detail.get('asset_type') input_block['item_index'] = detail.get('item_index') # Get correct asset field if input_block['asset_type'] == 'cover': input_block['asset_field'] = 'artwork' if input_block['asset_type'] == 'audio': input_block['asset_field'] = 'track' input_row = { 'item_index': input_block['item_index'], 'asset_type': input_block['asset_type'], 'correlation_id': input_block['correlation_id'], 'product_id': input_block['product_detail'].get('product_id'), 'upc': input_block['product_detail'].get('upc'), 'vendor_id': input_block['product_detail'].get('vendor_id'), 'subaccount_id': input_block['product_detail'].get('subaccount_id'), 'release_name': input_block['product_detail'].get('release_name'), 'project_code': input_block['product_detail'].get('project_code'), 'project_name': input_block['product_detail'].get('project_name'), 'project_id': input_block['product_detail'].get('project_id'), 'bucket': input_block['product_detail'].get(input_block['asset_field'], {}).get('bucket') if input_block['asset_field'] else None, 'key': input_block['product_detail'].get(input_block['asset_field'], {}).get('key') if input_block['asset_field'] else None, 'filename': input_block['product_detail'].get( input_block['asset_field'], {}).get('filename') if input_block[ 'asset_field'] else None, 'ows_assets_filename': input_block['product_detail'].get( input_block['asset_field'], {}).get( 'ows_assets_filename') if input_block['asset_field'] else None, } return input_row if __name__ == '__main__': print(f'SFN Execution ARN: {SFN_EXEC_ARN}') print('Getting Execution History.... This will take some time.') exec_list = sfn.get_execution_log_by_arn(SFN_EXEC_ARN, max_res=200) print('Converting history to dict') output = process_exec_list(exec_list) # Convert to dataframe print('Converting to Pandas DataFrame') df = pandas.DataFrame(output) # Write to Excel print(f'Writing XLSX out to: {SFN_EXEC_NAME}.xlsx') df.to_excel(SFN_EXEC_ARN+'.xlsx')