import os import sentry_sdk from config import config from connection import snowflake from queries.mysql import fetch_contributions_for_jobs, \ filter_contributions_in_queue from queries.mysql import purge_contributions_from_queue from queries.queries import get_unchanged_contribution log = config.get_logger() sentry_sdk.init(config.SENTRY_DSN) def purge_over_queued_contributions(cmm_id, job_last_updated, job_status): BATCH_SIZE = 500 log.info( f'Starting for store: {cmm_id} for date < {job_last_updated} ' f'and status: {job_status}' ) has_results = True purged_count = 0 try: sf_cursor = snowflake.get_cursor(target="contributions") offset = 0 while has_results: log.info(f'Getting contributions with offset: {offset}') job_result, rowcount = fetch_contributions_for_jobs( cmm_id, job_last_updated, job_status, BATCH_SIZE, offset ) has_results = bool(job_result) log.info(f'Contributions rowcount: {rowcount}, Ids: {job_result}') if not has_results: break offset += BATCH_SIZE queue_result, queue_count = filter_contributions_in_queue(job_result, cmm_id) log.info(f'In Queue rowcount: {queue_count}, Ids: {queue_result}') if queue_count < 1: continue log.info('Checking when these contributions were last updated.') unchanged_ids = get_unchanged_contribution( sf_cursor, job_last_updated, queue_result ) if len(unchanged_ids) > 0: log.info(f'for store: {cmm_id} Purging rowcount: {len(unchanged_ids)}, Ids: {unchanged_ids}') purge_contributions_from_queue(unchanged_ids, cmm_id) purged_count += len(unchanged_ids) log.info(f'Total {purged_count} records purged.') except Exception as e: log.error(e) raise e finally: log.info("Finished purge_over_queued_contributions") snowflake.close() if __name__ == "__main__": """Main entrypoint function.""" cmm_id = os.environ.get("CMM_ID") job_last_updated = os.environ.get("JOB_UPDATED_DATE") job_status = os.environ.get("JOB_STATUS", "complete") if not cmm_id or not job_last_updated: log.info('Missing require env variables: CMM_ID & JOB_UPDATED_DATE.') log.info('Run cmd: python cleanup_delivery_queue.py') exit(1) purge_over_queued_contributions(cmm_id, job_last_updated, job_status)