"""Queries.""" from lambdacommon import util _SQL_SET_JOB_STATUS = """UPDATE encoding_queue_detail SET `status` = %s, `error_log` = %s WHERE encoding_queue_detail_id = %s AND `status` = %s""" _SQL_GET_JOB = """SELECT * FROM encoding_queue_detail WHERE encoding_queue_detail_id = %s""" _SQL_GET_ACTIVE_DELIVERY_STORES = """ WITH counts AS ( SELECT eqd.dms_master_master_id, eq.encoding_order_type, COUNT(IF(eqd.status = 'new', 1, NULL)) AS count_new, COUNT(IF(eqd.status = 'encoded', 1, NULL)) AS count_encoded, COUNT(IF(eqd.status = 'ready_to_deliver', 1, NULL)) AS count_ready_to_deliver FROM encoding_queue_detail eqd JOIN encoding_queue eq ON eq.encoding_queue_id = eqd.encoding_queue_id WHERE eqd.status IN ( 'new', 'encoded', 'ready_to_deliver' ) GROUP BY 1, 2 ) SELECT dds.dms_master_master_id, dds.order_type, dds.batch_delivery, dds.encoding, dds.third_party_credential_set_id, counts.count_new, counts.count_encoded, counts.count_ready_to_deliver FROM counts JOIN dms_delivery_spec dds ON dds.dms_master_master_id = counts.dms_master_master_id AND dds.order_type = counts.encoding_order_type AND dds.delivery = 'Y' ; """ _SQL_GET_BATCH = """SELECT * FROM delivery_batch WHERE delivery_batch_id = %s""" _SQL_SET_BATCH_STATUS = """UPDATE delivery_batch SET `status` = %s WHERE delivery_batch_id = %s AND `status` = %s""" _SQL_GET_PRODUCT_ID_BY_UPC = """SELECT release_id, upc FROM releases WHERE upc IN (%s)""" _SQL_GET_TRACK_COUNT_IN_UPCS = """SELECT count(t.id) as total FROM art_relations.track t WHERE upc IN ({}) """ _SQL_GET_STORE_INFO = """SELECT customer_name, ddex_party_id FROM customer_master_master WHERE customer_master_master_id = %s""" _SQL_GET_STORAGE_DRIVES = """ SELECT sd.storage_drive_id, sd.drive_location, s.storage_id, s.unix_mount, INET_NTOA(s.local_ip) as local_ip, INET_NTOA(s.remote_ip) as remote_ip FROM storage_drive sd INNER JOIN storage_drive_asset_folder sdaf ON sd.storage_drive_id = sdaf.storage_drive_id INNER JOIN storage s ON s.storage_id = sd.storage_id WHERE s.physical_location_id = %s AND sdaf.asset_type_id = %s """ # ddtrace's "TracedCursor" class wraps callproc(), and # identifies 'args' as a positional arg rather than a keyword # arg. However, the PyMySQL spec lists 'args' as a keyword arg. # This is a small bug with ddtrace + PyMySQL + callproc(), so as a # workaround 'args' is passed into callproc() as a positional arg. def _callproc(cursor, procname, args=()): """Shell for callproc() function to gracefully handle 'args' bug. Args: cursor (Cursor): PyMySQL Cursor object procname (str): Name of procedure to execute. args (tuple): Parameters to use with procedure. """ cursor.callproc(procname, args) def set_job_status( from_status, to_status, eqd_id, error_log=None, conn_info=None): """Set the job status. Args: from_status (str): A valid from status to_status (str): A valid to status eqd_id (id): A valid encoding_queue_detail_id value error_log (str): An optional error log string conn_info (dict): An optional connection info dict Returns: bool: Affected rows. """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: success = cursor.execute( _SQL_SET_JOB_STATUS, (to_status, error_log, eqd_id, from_status)) conn.commit() return success def get_job(eqd_id, conn_info=None): """Get the job. Args: eqd_id (int): A valid encoding_queue_detail_id value conn_info (dict): An optional connection info dict Returns: dict: The status key with value. """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: cursor.execute( _SQL_GET_JOB, (eqd_id,)) return cursor.fetchone() def set_job_to_queued_for_delivery(eqd_id, conn_info=None): """Set job to queue_for_delivery helper. Args: eqd_id (int): A valid encoding_queue_detail_id value conn_info (dict): An optional connection info dict Returns: bool: Affected rows. """ return set_job_status( 'ready_to_deliver', 'queued_for_delivery', eqd_id, conn_info=conn_info) def get_connection_limit(dms_id, order_type, conn_info=None): """Get connection limit. Args: dms_id (int): A valid dms id order_type (str): A valid order type conn_info (dict): An optional connection info dict Returns: dict: A single result from the SP. """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc(cursor, 'sp_get_connection_limit', (dms_id, order_type)) return cursor.fetchone() def get_dms_delivery_spec(dms_id, order_type, conn_info=None): """Get dms delivery spec. Args: dms_id (int): A valid dms id order_type (str): A valid order type conn_info (dict): An optional connection info dict Returns: dict: A single result from the SP. """ return get_dms_delivery_specs(dms_id, order_type, conn_info, True) def get_dms_delivery_specs(dms_id, order_type, conn_info=None, fetch_one=False): """Get dms delivery spec. Args: dms_id (int): A valid dms id order_type (str): A valid order type conn_info (dict): An optional connection info dict fetch_one (bool): Flag for fetching only one record Returns: dict: A single result from the SP. """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_dms_delivery_spec', (dms_id, order_type)) if fetch_one: return cursor.fetchone() return cursor.fetchall() def get_jobs_for_batch_delivery(dms_id, allowed_encoders, conn_info=None): """Get jobs for batch delivery. Args: dms_id (int): A valid dms id allowed_encoders (list): A list of encoder ids conn_info (dict): An optional connection info dict Returns: list(dict) or tuple: A list of jobs """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: allowed_encoders = ','.join(map(str, allowed_encoders)) _callproc( cursor, 'sp_get_jobs_for_batch_delivery', (dms_id, allowed_encoders)) return cursor.fetchall() def get_jobs_for_nonbatch_delivery(dms_id, limit_to, allowed_encoders, conn_info=None): """Get jobs for non-batch delivery. Args: dms_id (int): A valid dms id limit_to (int): The sql limit in rows allowed_encoders (list): A list of encoder ids conn_info (dict): An optional connection info dict Returns: list(dict) or tuple: A list of jobs """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: allowed_encoders = ','.join(map(str, allowed_encoders)) _callproc( cursor, 'sp_get_jobs_for_nonbatch_delivery', (dms_id, limit_to, allowed_encoders)) return cursor.fetchall() def get_jobs_for_delivery( dms_id, limit_to, allowed_encoders, is_batch, conn_info=None): """Get jobs for batch or non-batch delivery. Args: dms_id (int): A valid dms id limit_to (int): The sql limit in rows allowed_encoders (list): A list of encoder ids is_batch (int): 1 or 0 for if its a batch job conn_info (dict): An optional connection info dict Returns: list(dict) or tuple: A list of jobs """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: allowed_encoders = ','.join(map(str, allowed_encoders)) _callproc( cursor, 'sp_get_jobs_for_delivery', (dms_id, limit_to, allowed_encoders, is_batch)) return cursor.fetchall() def get_active_delivery_stores(conn_info=None): """Get the active delivery stores. Args: conn_info (dict): An optional connection info dict Returns: list(dict): Returns list of dms ids and their order types info """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: cursor.execute(_SQL_GET_ACTIVE_DELIVERY_STORES) results = cursor.fetchall() active_stores = [] for store_id in set(x['dms_master_master_id'] for x in results): order_types = list(map( lambda x: { 'order_type': x['order_type'], 'batch_delivery': x['batch_delivery'], 'encoding': x['encoding'], 'third_party_credential_set_id': x['third_party_credential_set_id'], 'count_new': x['count_new'], 'count_encoded': x['count_encoded'], 'count_ready_to_deliver': x['count_ready_to_deliver'] }, filter( lambda x: x['dms_master_master_id'] == store_id, results))) active_stores.append({ 'dms_master_master_id': store_id, 'order_types': order_types}) return active_stores def get_renew_delivery_jobs(dms_id, limit_to, allowed_encoders, conn_info=None): """Get delivery jobs for Renew. Args: dms_id (int): A valid dms id limit_to (int): The sql limit in rows allowed_encoders (list): A list of encoder ids conn_info (dict): An optional connection info dict Returns: list(dict): Returns list of dms ids and their encoder ids """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: allowed_encoders = ','.join(map(str, allowed_encoders)) _callproc( cursor, 'sp_get_renewed_dd_items', (dms_id, allowed_encoders, limit_to)) results = cursor.fetchall() if type(results) is dict: return [results] return results def get_batch(batch_id, conn_info=None): """Get the batch. Args: batch_id (int): A valid delivery_batch_id value conn_info (dict): An optional connection info dict Returns: dict: Dictionary representation of row in delivery_batch. """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: cursor.execute( _SQL_GET_BATCH, (batch_id,)) return cursor.fetchone() def get_batches_to_close(conn_info=None): """Get ready_to_close batches for find_batch_to_close. Args: conn_info (dict): An optional connection info dict Returns: list(dict): Returns list of batch ids to close """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: cursor.callproc( 'sp_get_next_batch_to_close') results = cursor.fetchall() if type(results) is dict: return [results] return results def update_batch_status( from_status, to_status, batch_id, conn_info=None): """Get status of a batch by its id. Args: from_status (str): Batches with this status to update to_status (str): What status to update to batch_id (int) What batch_id to update conn_info (dict): An optional connection info dict Returns: bool: Affected rows. """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: success = cursor.execute( _SQL_SET_BATCH_STATUS, (to_status, batch_id, from_status)) return success def get_product_ids_by_upcs(upcs, chunk_size=500, conn_info=None): """Get release_ids by upcs. Args: upcs (list): A list of UPC's chunk_size (int): Size of each chunk to sent to query conn_info (dict): An optional connection info dict Returns: list: A list of dicts of release_id and upc. """ data = list() chunks = [upcs[i:i + chunk_size] for i in range(0, len(upcs), chunk_size)] with util.ar_connection(conn_info) as conn: with conn.cursor() as cursor: for chunk in chunks: format_string = ','.join(['%s'] * len(chunk)) cursor.execute( _SQL_GET_PRODUCT_ID_BY_UPC % format_string, tuple(chunk)) data += cursor.fetchall() return data def get_open_batch_count(dms_id, order_type='release', conn_info=None): """Get number of open batches for dms_id. Args: dms_id (int): Store ID order_type (string): Order type conn_info (dict): An optional connection info dict Returns: dict: Open batch count, number of open jobs and highest job priority """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_get_open_batch_count', (dms_id, order_type)) return cursor.fetchone() def get_encoded_jobs_info(dms_id, encoder_ids, order_type, conn_info=None): """Get encoded jobs info per order type by dms id. Args: dms_id (int): Store ID encoder_ids (list(int)): List of encoder ids to get jobs for order_type (string): Order type conn_info (dict): An optional connection info dict Returns: dict or None: A dict containing a store id, count of encoded jobs and highest job priority """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_get_encoded_jobs_info', ( dms_id, ','.join(map(str, encoder_ids)), order_type ) ) return cursor.fetchone() def create_delivery_batch( dms_id, remote_dir, order_type, number_per_batch, conn_info=None): """Create a delivery batch. Args: dms_id (int): Store ID remote_dir (string): Remote directory path string order_type (string): Order Type number_per_batch (int): Number of jobs per batch conn_info (dict): An optional connection info dict Returns: batch_id (int or None): Batch id of newly created batch """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_create_delivery_batch', (dms_id, remote_dir, order_type, number_per_batch)) results = cursor.fetchone() return results['batch_id'] if results['batch_id'] else None def set_job_for_delivery(dms_id, order_type, conn_info=None): """Set all encoded jobs to ready_to_deliver. Args: dms_id (int): Store ID order_type (string): Order type conn_info (dict): An optional connection info dict """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_set_job_for_delivery', (dms_id, order_type)) return True def close_delivery_batch(batch_id, conn_info=None): """Sets batch status to closed. Args: batch_id (int): Delivery Batch ID conn_info (dict): An optional connection info dict """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_close_delivery_batch', (batch_id,)) return True def get_batch_info(batch_id, conn_info=None): """Get batch information. Args: batch_id (int): Delivery Batch ID conn_info (dict): An optional connection info dict Returns: results (list or None): Batch info """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_get_batch_info', (batch_id,)) return cursor.fetchall() def get_release_priority_info( substore_ids, local_focus_territories, upcs, conn_info=None): """Get release priority info. Args: substore_ids (str): Sub-store ID's local_focus_territories (str): Territories (country ID's) upcs (str): List of UPC's conn_info (dict): An optional connection info dict Returns: results (list or None): Release priority info """ with util.ar_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_get_release_priority_info', (substore_ids, local_focus_territories, upcs)) return cursor.fetchall() def get_dms_master_contact_info(dms_id, conn_info=None): """Get DMS master contact info. Args: dms_id (int): Customer_master_master.customer_master_master_id in art_relations conn_info (dict): An optional connection info dict Returns: results (list or None): DMS master contact info """ with util.ar_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_get_dms_master_contact_info', (dms_id,)) return cursor.fetchall() def get_customer_name(dms_contact_id, conn_info=None): """Get DMS master name. Args: dms_contact_id (int): customer_master_master_contact.customer_master_master_contact_id in art_relations conn_info (dict): An optional connection info dict Returns: results (dict or Tuple): DMS master name """ with util.ar_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_get_customer_name', (dms_contact_id,)) return cursor.fetchone() def get_local_focus_territory(dms_id, conn_info=None): """Get local focus territory. Args: dms_id (int): customer_master_master_contact.customer_master_master_contact_id in art_relations conn_info (dict): An optional connection info dict Returns: results (dict or None): Local focus territory id. """ with util.ar_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_get_local_focus_territory', (dms_id,)) return cursor.fetchone() def get_customer_id(dms_id, conn_info=None): """Get customer id. Args: dms_id (int): customer_master_master_contact.customer_master_master_contact_id in art_relations conn_info (dict): An optional connection info dict Returns: results (list or empty tuple): Customer ids. """ with util.ar_connection(conn_info) as conn: with conn.cursor() as cursor: _callproc( cursor, 'sp_get_customer_id', (dms_id,)) return cursor.fetchall() def get_track_count_in_upcs(upcs, conn_info=None): """Get the number of tracks that belong in a list of upcs. Args: upcs (list): List of upcs. conn_info (dict): An optional connection info dict Returns: results (dict): Track count. """ with util.ar_connection(conn_info) as conn: with conn.cursor() as cursor: cursor.execute( _SQL_GET_TRACK_COUNT_IN_UPCS.format( ','.join([str(upc) for upc in upcs]))) return cursor.fetchone() def get_release_artist_genre_info(dms_ids, upc, conn_info=None): """Get release artist genre info. Args: dms_ids (list): List of DMS customer ids. upc (int): UPC number. conn_info (dict): An optional connection info dict. Returns: results (dict or None): Release artist genre info. """ with util.ar_connection(conn_info) as conn: with conn.cursor() as cursor: dms_customer_ids = ','.join(str(dms_id) for dms_id in dms_ids) _callproc( cursor, 'sp_get_release_artist_genre_info', (dms_customer_ids, upc)) return cursor.fetchone() def get_release_track_info(upc_list, include_track_code=True, conn_info=None): """Get basic release and track metadata from art_relations. Args: upc_list (list): A list of UPC's include_track_code (bool): Indicates whether to include track codes conn_info (dict): An optional connection info dict. Returns: results (list): List of release/track metadata """ with util.ar_connection(conn_info) as conn: with conn.cursor() as cursor: upcs = ','.join(str(upc) for upc in upc_list) _callproc( cursor, 'sp_get_release_track_info', (upcs, 1 if include_track_code else 0)) return cursor.fetchall() def get_store_info(dms_id, conn_info=None): """Get store name and DDEX party ID. Can be extended to get more data. Args: dms_id (int): Store ID from customer_master_master table in AR conn_info (dict): An optional db connection info dict. Returns: results (dict or None): Store record that contains DDEX party ID """ with util.ar_connection(conn_info) as conn: with conn.cursor() as cursor: cursor.execute(_SQL_GET_STORE_INFO, (dms_id,)) return cursor.fetchone() def get_storage_drives(physical_location_id, asset_type_id, conn_info=None): """Get a list of storage drives for the given physical location and asset type. Args: physical_location_id (int): Physical location ID asset_type_id (int): Asset type ID conn_info (dist): An optional db connection info dict. Returns: results (list): List of storage drives that are for given asset type """ with util.dd_connection(conn_info) as conn: with conn.cursor() as cursor: cursor.execute( _SQL_GET_STORAGE_DRIVES, (physical_location_id, asset_type_id)) return cursor.fetchall()