"""Queries for mysql connector.""" from connection.mysql import RdsHelper def fetch_contributions_for_jobs(cmm_id, job_last_updated, job_status, limit, offset): """Contributions for jobs that were completed before job_last_updated""" query = ( f"""SELECT c.contribution_uuid FROM claiming_order co INNER JOIN claiming_job cj ON co.id = cj.order_id INNER JOIN claims c ON cj.id = c.job_id WHERE co.customer_master_master_id = {cmm_id} AND cj.status = '{job_status}' AND cj.last_updated < '{job_last_updated}' ORDER BY c.contribution_uuid LIMIT {limit} OFFSET {offset} """ ) result, rowcount = RdsHelper.execute_rowcount_query(query) ids = [x[0] for x in result] return ids, rowcount def filter_contributions_in_queue(contribution_ids, cmm_id): """Filter contributions that are in queue.""" ids_string = "','".join(contribution_ids) query = f"""SELECT dq.contribution_uuid FROM delivery_queue dq WHERE dq.contribution_uuid in ('{ids_string}') AND dq.customer_master_master_id = {cmm_id} """ result, rowcount = RdsHelper.execute_rowcount_query(query) ids = [x[0] for x in result] return ids, rowcount def purge_contributions_from_queue(contribution_ids, cmm_id): ids_string = "','".join(contribution_ids) delete = f"""delete from delivery_queue where contribution_uuid in ('{ids_string}') and customer_master_master_id = {cmm_id}""" return RdsHelper.execute_delete_query(delete) def get_error_jobs_with_queued_claims(limit, offset): """Select job and order id.""" query = f"""SELECT c.job_id, cj.order_id FROM claims c INNER JOIN claiming_job cj ON cj.id = c.job_id AND cj.status = 'error' WHERE c.status IN ('queued','started') GROUP BY c.job_id LIMIT {limit} OFFSET {offset} """ result, rowcount = RdsHelper.execute_rowcount_query(query) ids = [{'job_id': x[0], 'order_id': x[1]} for x in result] return ids, rowcount def update_claim_status_for_job(job_id, status): """Update claim status. job_id: job unique identifier status: status to update Returns: int: with the number of rows updated """ query = f"""UPDATE claims SET status = '{status}' WHERE job_id = {job_id} """ params = (status, job_id) RdsHelper.check_connection() cursor = RdsHelper.cnx.cursor() cursor.execute(query, params=params) RdsHelper.cnx.commit() result = cursor.rowcount cursor.close() return result def update_claims_status(claims, status): RdsHelper.check_connection() cursor = RdsHelper.cnx.cursor() query = """UPDATE claims SET status = %s, message = %s, contribution_version = %s WHERE id = %s """ for each_claim in claims: params = ( status, each_claim['reason'], each_claim['versionId'], each_claim['claimId'] ) cursor.execute(query, params=params) RdsHelper.cnx.commit() result = cursor.rowcount cursor.close() return result