from config import config import mysql.connector MIGRATION_ID = "7c86a224-053c-4b4d-a170-78de2f91fe75" MIGRATION_DATE = "2021-01-01" class RdsHelper: cnx = mysql.connector.connect( user=config.NR_DELIVERY_USERNAME, password=config.NR_DELIVERY_PASSWORD, host=config.NR_DELIVERY_HOSTNAME, database=config.NR_DELIVERY_DB_NAME, ) @staticmethod def renew_connection(): RdsHelper.cnx.cmd_reset_connection() if RdsHelper.cnx.is_connected(): return else: RdsHelper.cnx.reconnect(attempts=3, delay=1) if not RdsHelper.cnx.is_connected(): try: RdsHelper.cnx.close() RdsHelper.cnx = mysql.connector.connect( user=config.NR_DELIVERY_USERNAME, password=config.NR_DELIVERY_PASSWORD, host=config.NR_DELIVERY_HOSTNAME, database=config.NR_DELIVERY_DB_NAME, ) except Exception as err: print(err) raise err @staticmethod def check_connection(): if not RdsHelper.cnx.is_connected(): RdsHelper.renew_connection() @staticmethod def execute_query(query): RdsHelper.check_connection() cursor = RdsHelper.cnx.cursor() cursor.execute(query) RdsHelper.cnx.commit() last_row_id = cursor.lastrowid cursor.close() return last_row_id @staticmethod def execute_rowcount_query(query): RdsHelper.check_connection() cursor = RdsHelper.cnx.cursor() cursor.execute(query) result = cursor.fetchall() row_count = cursor.rowcount cursor.close() return result, row_count @staticmethod def execute_delete_query(query): RdsHelper.check_connection() cursor = RdsHelper.cnx.cursor() cursor.execute(query) RdsHelper.cnx.commit() result = cursor.rowcount cursor.close() return result def commit(): RdsHelper.cnx.commit() def create_order(store_id): query = f"INSERT INTO claiming_order(customer_master_master_id, created_by, created_at, last_updated, status)" values = f" values({store_id}, '{MIGRATION_ID}', '{MIGRATION_DATE}', '{MIGRATION_DATE}', 'success')" return RdsHelper.execute_query(query + values) def insert_contributor_society_into_jobs_table(order_id, contributor_uuid): query = "insert into claiming_job(order_id, contributor_uuid, status, last_updated, last_updated_by) " values = f"values({order_id}, '{contributor_uuid}', 'complete', '{MIGRATION_DATE}', '{MIGRATION_ID}')" return RdsHelper.execute_query(query + values) def insert_job_contribution_into_claims_table( job_id, contribution_uuid, status, creation_date ): insert = ( "insert into claims(contribution_uuid, job_id, status, last_updated, message ) " ) values = f"values('{contribution_uuid}', {job_id}, '{status}', '{creation_date}', 'migration');" query = insert + values return RdsHelper.execute_query(query) def add_to_queue(contribution, cmm_id, contributor): insert = f"insert into delivery_queue (contribution_uuid, customer_master_master_id, contributor_uuid) values ('{contribution}', '{cmm_id}', '{contributor}');" try: return RdsHelper.execute_query(insert) except mysql.connector.errors.IntegrityError as e: return None def add_many_to_queue(values): insert = f"insert ignore into delivery_queue (contribution_uuid, customer_master_master_id, contributor_uuid) values " insert += values return RdsHelper.execute_query(insert) def fetch_excluded_contributions(): query = ("select claims.contribution_uuid, orders.customer_master_master_id from claims " + "inner join claiming_job as jobs on claims.job_id = jobs.id " + "inner join claiming_order as orders on jobs.order_id = orders.id " + "inner join delivery_queue as dq on dq.contribution_uuid = claims.contribution_uuid and dq.customer_master_master_id = orders.customer_master_master_id " + "where claims.status = 'excluded';") return RdsHelper.execute_rowcount_query(query) def fetch_claims_for_contribution_cmo(contribution_uuid, cmm_id): query = ("select claims.contribution_uuid, orders.customer_master_master_id, claims.last_updated, claims.status, claims.job_id from claims " + "inner join claiming_job as jobs on claims.job_id = jobs.id " + "inner join claiming_order as orders on jobs.order_id = orders.id " + f"where claims.contribution_uuid = '{contribution_uuid}' AND orders.customer_master_master_id = {cmm_id} " + "order by claims.last_updated desc;") return RdsHelper.execute_rowcount_query(query) def delete_excluded_contributions_from_queue(contribution_uuid, cmm_id): delete = (f"delete from delivery_queue where contribution_uuid = '{contribution_uuid}' and customer_master_master_id = {cmm_id}") return RdsHelper.execute_delete_query(delete) def close(): if RdsHelper.cnx: RdsHelper.cnx.close()