""" script to prepare for download of some additional BUs To run: 1. make venv 2. . venv/bin/activate 3. python add_new_bu.py --env=dev --start-date=20210731 --report-name=activityfeed ( python add_new_bu.py -h) """ import argparse import logging import time from datetime import datetime, timedelta, timezone from sqlalchemy.orm.attributes import flag_modified import boto3 from db_schema.schemas import slz from slz_storage.postgres import connection FIRST_AVAILABLE_DATE_ACTIVITYFEED = { 'sme': { 'US_ThreeSixZero': '2021-02-05 19:57:32' } } FIRST_AVAILABLE_DATE_MEMBERSLOGIN = { 'sme': { 'AT_MediaSendsNoDOI': '2019-07-30 07:22:42', 'DE_PromotionNoDOI': '2019-07-30 07:22:42', 'US_ThreeSixZero': '2019-07-27 15:29:03' } } FIRST_AVAILABLE_DATE_MEMBERSEXTENDED = {} REPORTS_MAPPER = { 'activityfeed': FIRST_AVAILABLE_DATE_ACTIVITYFEED, 'memberslogin': FIRST_AVAILABLE_DATE_MEMBERSLOGIN, 'membersextended': FIRST_AVAILABLE_DATE_MEMBERSEXTENDED, } UOW_MOCK = "appreciationengine-{day}-{licensor_name}-{report_name}-v1" def run(params: argparse.Namespace): logging.basicConfig(format='[CREATE UOW/CS] %(levelname)s:%(message)s', level=logging.INFO) logging.info("Start script execution") secretsmanager = boto3.client('secretsmanager') pg = connection.get_session_from_secret_key(secretsmanager, f'delphi/{params.env}/slz/storage/pg_proxy/user') try: first_available_dates = REPORTS_MAPPER[params.report_name] start_date = params.start_date td = timedelta(days=1) now = datetime.now(tz=timezone.utc) + timedelta(hours=1) report: slz.Report = pg.query(slz.Report).filter( slz.Report.report_name == params.report_name).one() licensors = { item.licensor_name: item.licensor_id for item in pg.query(slz.Licensor).all() } while True: processed = [] for licensor_name, dates in first_available_dates.items(): unit_of_work_code = UOW_MOCK.format( day=start_date.strftime('%Y%m%d'), licensor_name=licensor_name, report_name=params.report_name ) # looking for existing UOW uow_record: slz.UnitOfWork = pg.query(slz.UnitOfWork).filter( slz.UnitOfWork.unit_of_work_code == unit_of_work_code ).order_by(slz.UnitOfWork.unit_of_work_id.desc()).limit(1).one_or_none() if uow_record: logging.info("Found uow record %s", unit_of_work_code) # if found, make it active again uow_record.completeness_status = slz.CompletenessStatusEnum.ACTIVE.value uow_record.created_at = now uow_record.last_updated_at = now uow_record.next_run_at = now + timedelta(seconds=10) uow_record.is_force_complete = False else: logging.info("Create uow record %s", unit_of_work_code) uow_record = slz.UnitOfWork( unit_of_work_code=unit_of_work_code, reprocess_id='', report_date=start_date.date(), report_id=report.report_id, licensor_id=licensors[licensor_name], version='v1', activity_status=slz.ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=slz.CompletenessStatusEnum.ACTIVE.value, unit_of_work_type=slz.UnitOfWorkTypeEnum.BACKFILL.value, next_run_at=now + timedelta(seconds=10), created_at=now, last_updated_at=now, is_force_complete=False, priority=slz.UnitOfWorkPriorityEnum.PRIORITY_6.value ) pg.add(uow_record) for context, earliest_date_str in dates.items(): if not earliest_date_str: continue earliest_date = datetime.strptime(earliest_date_str, '%Y-%m-%d %H:%M:%S') if start_date.date() >= earliest_date.date(): # try to find existing content_status_record content_status_record: slz.ContentStatus = pg.query(slz.ContentStatus).filter( slz.ContentStatus.unit_of_work_id == uow_record.unit_of_work_id, slz.ContentStatus.context == context ).one_or_none() if content_status_record: logging.info("Found cs %s, %s", unit_of_work_code, context) content_status_record.content_status = slz.ContentStatusEnum.MISSING content_status_record.meta_data = {} flag_modified(content_status_record, "meta_data") else: logging.info("create cs %s, %s", unit_of_work_code, context) content_status_record = slz.ContentStatus( content_name=f"{context}_{start_date.strftime('%Y%m%d')}.json", context=context, content_status=slz.ContentStatusEnum.MISSING, failure_count=0, latest_job_id='backfill', created_at=now, last_checked_at=now, sub_content=None, metadata_process_status=slz.ContentMetadataStatusEnum.COMPLETED, ) content_status_record.unit_of_work = uow_record pg.add(content_status_record) processed.append(content_status_record) if not processed: break try: pg.commit() except Exception as err: logging.error("Error during commit: %s", str(err)) pg.rollback() start_date -= td now += timedelta(minutes=1) time.sleep(0.5) finally: pg.close() logging.info("Finish script execution") if __name__ == '__main__': def date_type(value): try: value = datetime.strptime(value, '%Y%m%d') except Exception: raise argparse.ArgumentTypeError("Not a valid date") return value parser = argparse.ArgumentParser() parser.add_argument("--env", type=str, required=True, choices=['dev', 'qa', 'stage', 'prod']) parser.add_argument("--report-name", type=str, required=True, choices=['memberslogin', 'activityfeed', 'membersextended']) parser.add_argument("--start-date", type=date_type, required=True, help="in format YYYYMMDD" ) params = parser.parse_args() run(params)