"""Hit ows-delivery-metadata and report HTTP response.""" import argparse import asyncio import csv import json import httpx """ SELECT eqd.upc, IFF(eq.META_UPDATE = 'Y', 'metadata_update', 'complete_album') AS DELIVERY_TYPE, MAX(eqd.DELIVERY_ENDED) AS dt FROM ORCHARD_APP_REPORTING_V2.REPORTSDD_DIRECT_DELIVERY.ENCODING_QUEUE_DETAIL eqd JOIN ORCHARD_APP_REPORTING_V2.REPORTSDD_DIRECT_DELIVERY.ENCODING_QUEUE eq ON eq.ENCODING_QUEUE_ID = eqd.ENCODING_QUEUE_ID JOIN ORCHARD_APP_REPORTING_V2.ART_RELATIONS_PROD_ART_RELATIONS.releases r ON r.upc = eqd.UPC AND r.deletions = 'N' AND r.release_status = 'in_content' WHERE eqd.DMS_MASTER_MASTER_ID = 286 AND eqd.STATUS = 'delivered' AND eqd.DELIVERY_ENDED > DATEADD(DAY, -102, CURRENT_DATE()) AND eqd.DELIVERY_ENDED < DATEADD(DAY, -2, CURRENT_DATE()) GROUP BY 1, 2 ORDER BY dt DESC ; """ async def main(): """Entrypoint.""" parser = argparse.ArgumentParser( formatter_class=argparse.ArgumentDefaultsHelpFormatter ) parser.add_argument( 'input_filename', type=str, help='Filename with UPCs as input.' ) parser.add_argument( 'dms_id', type=int, help='DMS ID of store to request data for.' ) parser.add_argument( '--env', type=str, choices=['prod', 'qa'], default='qa', help='Target application environment.' ) parser.add_argument( '--workers', type=int, default=20, help='Number of concurrent requests to make.' ) parser.add_argument( '--after', type=int, help='Process rows after this index only.' ) parser.add_argument( '--retries', type=int, default=3, help='Number of times to retry HTTP issues.' ) parser.add_argument( '--timeout', type=int, default=60, help='Network timeout.' ) args = parser.parse_args() print(args) # load UPCs from file with open(f'/tmp/input/{args.input_filename}', 'r') as f: reader = csv.DictReader(f) rows = [ (x['UPC'], x['DELIVERY_TYPE']) for x in reader ] # skip to point in file if args.after: rows = rows[args.after + 1:] # put UPCs into queue queue = asyncio.Queue() for row in rows: queue.put_nowait(row) # work through queue await asyncio.gather( *[ asyncio.create_task( process_queue( queue, args.env, args.dms_id, args.retries, args.timeout ) ) for _ in range(args.workers) ] ) async def process_queue(queue, env, dms_id, retries, timeout): """Process event from queue.""" while not queue.empty(): (upc, delivery_type) = await queue.get() print(f'start {upc} {delivery_type}') # make request url = f'https://{env}-ows-delivery-metadata.theorchard.io/products/{upc}/stores/{dms_id}/delivery_metadata' # noqa:E501 params = { 'delivery_type': delivery_type, 'version': '4.3', 'format': 'ddex_ern' } client_params = { 'transport': httpx.AsyncHTTPTransport(retries=retries), 'timeout': timeout } async with httpx.AsyncClient(**client_params) as client: response = await client.get( url, params=params ) # handle response details = None if response.status_code != 200: try: data = response.json() except json.decoder.JSONDecodeError: details = 'unable to parse JSON' else: details = data.get('code') if details == 'pydantic_validation_error': details += ':'.join([ x.get('msg') for x in data.get('errors', []) ]) print(response.status_code, upc, delivery_type, details) queue.task_done()