"""Lambda handler for proxying request from SNS to StepFunction.""" import os import json import argparse import datetime import boto3 import botocore from connectors import mysql from models import asset_upload from models import oat_asset_upload sfn_client = boto3.client('stepfunctions', region_name='us-east-1') s3_client = boto3.client('s3', region_name='us-east-1') workstation_config = { 'dev': { 'STATEMACHINE_ARN': 'arn:aws:states:us-east-1:103233932089:stateMachine:dev_assets_transcoding_v2', 'BUCKET': 'dev-orcd-raw-assets' }, 'qa': { 'STATEMACHINE_ARN': 'arn:aws:states:us-east-1:437795906767:stateMachine:qa_assets_transcoding_v2', 'BUCKET': 'qa-orcd-raw-assets' }, 'prod': { 'STATEMACHINE_ARN': 'arn:aws:states:us-east-1:437795906767:stateMachine:prod_assets_transcoding_v2', 'BUCKET': 'prod-orcd-raw-assets' } } podcast_config = { 'dev': { 'STATEMACHINE_ARN': 'arn:aws:states:us-east-1:103233932089:stateMachine:dev-asset-transcoder-state-machine', 'BUCKET': 'dev-orcd-asset-transcoder-input' }, 'qa': { 'STATEMACHINE_ARN': 'arn:aws:states:us-east-1:437795906767:stateMachine:qa-asset-transcoder-state-machine', 'BUCKET': 'qa-orcd-asset-transcoder-input' }, 'prod': { 'STATEMACHINE_ARN': 'arn:aws:states:us-east-1:437795906767:stateMachine:prod-asset-transcoder-state-machine', 'BUCKET': 'prod-orcd-asset-transcoder-input' } } def process_file(event, state_machine): """Step function entry point. Args: event (dict): Information about uploaded image. Returns: dict: Dict with start execution results """ try: response = sfn_client.start_execution( stateMachineArn=state_machine, input=json.dumps(event) ) return response except botocore.exceptions.ClientError as e: print(e) def get_missing_podcast_files(request_date): with mysql.oat_au_db_session() as session: uploads = oat_asset_upload.get_broken_uploads( session, request_date) return [upload.filename for upload in uploads] def get_missing_workstation_files(request_date): with mysql.au_db_session() as session: uploads = asset_upload.get_broken_uploads( session, request_date) return [upload.filename for upload in uploads] missing_keys = [] if __name__ == "__main__": argparser = argparse.ArgumentParser(prog='process assets') argparser.add_argument('-env', required=True, help='env', default='qa') argparser.add_argument('-invoke_type', required=True, help='type: workstation or podcast', default='workstation') argparser.add_argument('-request_date', required=True, help='YYYY-MM-DD') args = argparser.parse_args() env = args.env invoke_type = args.invoke_type request_date = datetime.datetime.strptime(args.request_date, '%Y-%m-%d') vars = workstation_config[env] if invoke_type == 'workstation' else podcast_config[env] if invoke_type == 'workstation': files = get_missing_workstation_files(request_date) else: files = get_missing_podcast_files(request_date) print('found {} files created since {}'.format(len(files), request_date)) print(vars) for file in files: response = s3_client.list_objects_v2( Bucket=vars['BUCKET'], Prefix=file ) if response["KeyCount"] == 1: key = response["Contents"][0]["Key"] ev = { "detail": { "requestParameters": { "bucketName": vars['BUCKET'], "key": key } } } r = process_file(ev, vars['STATEMACHINE_ARN']) print(r) else: missing_keys.append(file) print('{} missing keys out of {} files'.format(len(missing_keys), len(files))) print(missing_keys)