import argparse from concurrent import futures from functools import wraps import json import os import sys import threading import time import csv import boto3 import requests _DEFAULT_POOL = futures.ThreadPoolExecutor() def threadpool(f, executor=None): @wraps(f) def wrap(*args, **kwargs): return (executor or _DEFAULT_POOL).submit(f, *args, **kwargs) return wrap class ProgressPercentage(object): def __init__(self, filename): self._filename = filename self._size = float(os.path.getsize(filename)) self._seen_so_far = 0 self._lock = threading.Lock() def __call__(self, bytes_amount): # To simplify we'll assume this is hooked up # to a single filename. with self._lock: self._seen_so_far += bytes_amount percentage = (self._seen_so_far / self._size) * 100 print( "\r%s %s / %s (%.2f%%)" % ( self._filename, self._seen_so_far, self._size, percentage)) #sys.stdout.flush() def get_assets_host(env): return 'https://{}-ows-assets.theorchard.io'.format(env) def create_token(env, headers): token_response = requests.get( '{}/upload-token'.format(get_assets_host(env)), headers=headers ) return token_response.json() def get_full_filename(token_response_body, upload_file): return '{}.{}'.format( token_response_body.get('filename'), upload_file.get('asset_type') ) def upload(product, upload_file, headers, token_response_body): upload_creds = token_response_body.get('credentials') bucket_name = token_response_body.get('bucket') dest_file_name = get_full_filename(token_response_body, upload_file) file_to_upload = upload_file.get('file') client = boto3.client( 's3', aws_access_key_id=upload_creds.get('aws_access_key_id'), aws_secret_access_key=upload_creds.get('aws_secret_access_key'), aws_session_token=upload_creds.get('token') ) s3_metadata = { 'asset_type': upload_file.get('asset_type').upper(), 'product_id': product.get('product_id'), 'original_filename': upload_file.get('original_filename'), 'upc': product.get('upc'), 'track_unique_id': product.get('tuid') } print('dest file: {}'.format(dest_file_name)) print('s3 metadata: {}'.format(json.dumps(s3_metadata))) print('content type: {}'.format(upload_file.get('content_type'))) client.upload_file( file_to_upload, bucket_name, dest_file_name, Callback=ProgressPercentage(file_to_upload), ExtraArgs={ 'Metadata': s3_metadata, 'ContentType': upload_file.get('content_type') } ) def poll_v2(env, token_response_body, upload_file, headers): status_url = '{}/v2/asset/status/{}'.format( get_assets_host(env), get_full_filename(token_response_body, upload_file) ) poll_status(status_url, 'v2', 'encoding_completed', headers) def post_v1(env, product, upload_file, headers, token_response_body): dest_file_name = get_full_filename(token_response_body, upload_file) upload_creds = token_response_body.get('credentials') post_payload = { 'asset_type': upload_file.get('asset_type'), 'product_id': int(product.get('product_id')), 'original_filename': upload_file.get('original_filename'), 'upc': str(product.get('upc')), 'filename': dest_file_name, 'track_unique_id': int(product.get('tuid')), 'token': upload_creds.get('token') } print('v1 post payload: {}'.format(json.dumps(post_payload))) post_response = requests.post( '{}/asset'.format(get_assets_host(env)), headers=headers, data=json.dumps(post_payload) ) print(post_response.json()) def poll_v1(env, token_response_body, upload_file, headers): status_url = '{}/asset/status?filename={}'.format( get_assets_host(env), get_full_filename(token_response_body, upload_file) ) poll_status(status_url, 'v1', 'finished', headers) def poll_status(status_url, version, finished_status, headers): print(status_url) last_status = None keep_going = True while keep_going: print_body = True status_response = requests.get(status_url, headers=headers) if status_response.status_code == 200: status_data = status_response.json() current_status = status_data['status'] print('{} status: {}'.format(version, current_status)) if current_status == finished_status: keep_going = False if current_status == last_status: print_body = False last_status = current_status else: print('{} status response code: {}'.format( version, status_response.status_code)) if print_body: print('response body: {}'.format(status_response.text)) time.sleep(1) @threadpool def v1_time(env, product, upload_file, headers, token_response_body): start_time = time.time() post_v1(env, product, upload_file, headers, token_response_body) poll_v1(env, token_response_body, upload_file, headers) end_time = time.time() return end_time - start_time @threadpool def v2_time(env, upload_file, headers, token_response_body): start_time = time.time() poll_v2(env, token_response_body, upload_file, headers) end_time = time.time() return end_time - start_time def write_csv(v1_time, v2_time, content_type, file_size): csv_dict = {} csv_dict['content_type'] = content_type csv_dict['file_size'] = file_size csv_dict['v1'] = v1_time csv_dict['v2'] = v2_time diff = v1_time - v2_time csv_dict['diff'] = diff csv_dict['percentage'] = diff / v2_time with open('v1_v2_results.csv', 'a') as f: w = csv.DictWriter(f, csv_dict.keys()) w.writerow(csv_dict) def main(): parser = argparse.ArgumentParser(prog='ows-assets comparison tool') parser.add_argument('-p', '--product_id', required=True, help='product id to upload an asset to') parser.add_argument('-t', '--tuid', required=False, default='0', help='track id to upload an asset to') parser.add_argument('-u', '--upc', default=None, required=False, help='UPC of the product, optional') parser.add_argument('-g', '--grass', default=25824, help='Grass Account Id') parser.add_argument('-o', '--orchard', default='alw:43807', help='Orchard User Id') parser.add_argument('-f', '--file', required=True, help='File to upload') parser.add_argument('-c', '--content_type', required=True, help='Content Type to set the S3 metadata to') parser.add_argument('-e', '--env', required=False, help='Environment (qa or prod)', default='qa') args = parser.parse_args() env = args.env print('Executing against the {} environment'.format(env)) product = { 'product_id': args.product_id, 'upc': args.upc, 'tuid': args.tuid } file_size = os.path.getsize(args.file) _, file_extension = os.path.splitext(args.file) upload_file = { 'file': args.file, 'original_filename': os.path.basename(args.file), 'asset_type': file_extension.replace('.', ''), 'content_type': args.content_type } request_headers = { 'Grass-Account-Type': 'vendor', 'Grass-Account-Id': str(args.grass), 'Content-Type': 'application/json', 'Correlation-Id': 'abc123', 'Orchard-User-Id': args.orchard } token_response_body = create_token(env, request_headers) upload(product, upload_file, request_headers, token_response_body) v2 = v2_time(env, upload_file, request_headers, token_response_body) v1 = v1_time( env, product, upload_file, request_headers, token_response_body ) write_csv(v1.result(), v2.result(), args.content_type, file_size) sys.exit(0) if __name__ == "__main__": main()