import argparse import json import os import sys import time import uuid import requests from constants import bulk_assets from constants import asset_targets from models.asset import Audio from integration_scripts import logger from integration_scripts import s3_backoff_utils as s3_utils from integration_scripts import orch_utils import config MAKE_OWS_REQUESTS = False MAKE_S3_REQUESTS = False IMAGE_TYPE = 'image/tiff' AUDIO_TYPE = 'audio/wav' allowed_content_types = [IMAGE_TYPE, AUDIO_TYPE] def get_assets_host(env): """Returns the location of the assets host.""" return 'https://{}-ows-assets.theorchard.io'.format(env) def request_filenames(vendor_id, alw_id, token, count, env): """Request unique filenames from ows_assets.""" request_headers = get_bulk_request_headers(vendor_id, alw_id) body_data = get_bulk_request_data(token, count) filename_list = list() if MAKE_OWS_REQUESTS: post_response = requests.post( url='{}/upload-token-bulk'.format(get_assets_host(env)), headers=request_headers, data=body_data ) filename_list = post_response.json().get('filenames') logger.info(post_response.json()) return filename_list def ingest_one( bucket_name, product, upload_file, unique_file_name, headers, token_str): file_to_upload = upload_file.get('file_and_path') response = None 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'), 'release_status': bulk_assets.release_status, 'is_correction': bulk_assets.is_correction } logger.info('Copy file on S3.') if MAKE_S3_REQUESTS: response = s3_utils.copy_from_backoff( bucket_name, file_to_upload, unique_file_name, Metadata=s3_metadata, ContentType=upload_file.get('content_type') # Callback=ProgressPercentage(file_to_upload), # ExtraArgs={ # 'Metadata': s3_metadata, # 'ContentType': upload_file.get('content_type') # } ) if config.SEND_API_V1: logger.info('Sending assets v1 request.') post_v1( config.ENVIRONMENT, product, upload_file, headers, unique_file_name, token_str) return unique_file_name, response def post_v1(env, product, upload_file, headers, dest_file_name, token): 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': token } if MAKE_OWS_REQUESTS: post_response = requests.post( '{}/asset'.format(get_assets_host(env)), headers=headers, data=json.dumps(post_payload) ) logger.info(post_response.json()) def poll_v2(env, dest_file_name, headers): status_url = '{}/v2/asset/status/{}'.format( get_assets_host(env), dest_file_name ) if MAKE_OWS_REQUESTS: poll_status(status_url, 'v2', 'encoding_completed', headers) else: logger.info(status_url) def poll_v1(env, dest_file_name, headers): status_url = '{}/asset/status?filename={}'.format( get_assets_host(env), dest_file_name ) if MAKE_OWS_REQUESTS: poll_status(status_url, 'v1', 'finished', headers) else: logger.info(status_url) def poll_status(status_url, version, finished_status, headers): logger.info(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: logger.info('{} status response code: {}', version, status_response.status_code) if print_body: logger.info('response body: {}', status_response.text) time.sleep(1) def get_product_request(product_id, upc=None, tuid=0): product = { 'product_id': product_id, 'upc': upc, 'tuid': tuid } return product def get_upload_file_request(file, ext, content_type): upload_file = { 'file_and_path': file, 'original_filename': os.path.basename(file), 'asset_type': ext.replace('.', ''), 'content_type': content_type } return upload_file def get_request_headers(grass=25824, orchard='alw:43807'): request_headers = { 'Grass-Account-Type': 'vendor', 'Grass-Account-Id': str(grass), 'Content-Type': 'application/json', 'Correlation-Id': uuid.uuid1(), 'Orchard-User-Id': orchard } return request_headers def get_bulk_request_headers(grass=25824, orchard_id='alw:43807'): request_headers = { 'Grass-Account-Type': 'vendor', 'Grass-Account-Id': str(grass), 'Content-Type': 'application/json', 'Correlation-Id': uuid.uuid1(), 'Orchard-User-Id': orchard_id } return request_headers def get_bulk_request_data(token='bulk_upload', count=100): return { 'token': token, 'bulk_upload_count': count } def send_asset(unique_filename, bucket, grass_id, orchard_id, content_type, original_file, product_id, upc, tuid, token): product = get_product_request(product_id, upc, tuid) _, file_extension = os.path.splitext(original_file) upload_file = get_upload_file_request( original_file, file_extension, content_type) request_headers = get_request_headers(grass_id, orchard_id) dest_file_name, _ = ingest_one( bucket, product, upload_file, unique_filename, request_headers, token) 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_id', default=25824, help='Grass Account Id') parser.add_argument('-o', '--orchard_id', 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, default=AUDIO_TYPE, help='Content Type to set the S3 ' 'metadata to') parser.add_argument('-e', '--env', required=False, help='Environment (qa or prod)', default='qa') parser.add_argument('-b', '--bucket', required=True, help='Target Upload' 'Bucket') args = parser.parse_args() env = args.env token = 'TEST_TOKEN' if args.content_type not in allowed_content_types: sys.exit('Content-type must be one of \'{}\''.format( '\', \''.join(allowed_content_types))) logger.info('Executing against the {} environment', env) vendor_id = args.grass_id alw_id = orch_utils.get_user_id_alw(vendor_id) # Shape assets into chunks of MAX_REQUEST_PER_LOOP chunk_list = [ [Audio(), Audio(), Audio()], [Audio(), Audio(), Audio()] ] # Determine target bucket if env == 'prod': target = asset_targets.PROD_BUCKET else: target = asset_targets.QA_BUCKET # For each chunk of size MAX_REQUEST_PER_LOOP for chunk in chunk_list: for asset in chunk: filename_list = request_filenames( vendor_id, alw_id, token, config.MAX_REQUEST_PER_LOOP, env ) filename = filename_list.pop() send_asset( filename, target, asset.vendor_id, asset.alw_id, asset.content_type, asset.source_audio_file, asset.product_id, asset.orch_upc, asset.tuid, token ) sys.exit(0) if __name__ == "__main__": main()