""" script To run: 1. make venv 2. . venv/bin/activate 3. python generate_uow.py --env=dev --start-date=20190601 --report-name=activityfeed ( python generate_uow.py -h) """ import argparse import logging from datetime import datetime, timedelta, timezone import boto3 from db_schema.schemas import slz from slz_storage.postgres import connection FIRST_AVAILABLE_DATE_ACTIVITYFEED = { 'sme': { 'CenturyMediaRecords': '2018-01-01 00:33:27', 'AR': '2019-01-02 18:46:46', 'Asia': '2017-12-01 10:49:14', 'AU_NZ': '2018-01-11 19:35:51', 'AT': '2018-09-17 16:41:00', 'BE': '2017-12-01 09:56:32', 'BR': '2018-01-01 00:08:10', 'BR_FiltrStore': None, 'CA': '2018-01-01 02:46:48', 'CentralAmericaCaribbean': '2019-05-24 20:34:15', 'CL': '2018-03-06 17:49:02', 'CO': '2018-01-02 01:30:38', 'DK': '2018-05-29 15:17:28', 'US': '2018-12-19 19:27:45', 'FI': '2018-05-25 10:31:28', 'FR': '2019-10-15 15:05:43', 'DE': '2018-01-01 00:00:44', 'HK_LiquidState': None, 'IN': '2019-07-26 13:11:38', 'ID': '2019-01-23 08:40:19', 'IE': '2018-09-18 10:53:12', 'IT': '2018-01-01 00:00:26', 'MX': '2018-01-01 00:14:47', 'NL': '2019-01-01 01:55:25', 'NO': '2018-01-01 00:44:47', 'PE': '2019-12-02 23:04:34', 'PL': '2019-01-04 14:49:00', 'PT': '2018-10-17 17:46:37', 'ZA': '2018-04-04 09:14:50', 'ES': '2018-01-01 00:17:09', 'SE': '2018-04-12 11:49:34', 'CH': '2018-01-01 14:10:49', 'TR': '2018-11-19 13:06:51', 'UK': '2018-01-01 00:06:51', 'ClassicalInternational': None, 'WestAfrica_NG': None, 'US_Arista': '2019-02-06 22:33:51', 'US_Columbia': '2018-01-01 00:00:59', 'US_CRM_ProdSupport': '2018-01-02 15:45:57', 'US_Epic': '2018-01-01 00:00:00', 'US_GDB_Sales': '2018-09-17 14:25:09', 'US_LatinIberia': None, 'US_LatinUS': '2018-01-01 00:58:47', 'US_Legacy': '2018-01-01 00:21:53', 'US_Masterworks': '2018-03-22 00:49:01', 'US_MonumentRecords': '2018-01-01 01:40:48', 'US_Nashville': '2018-01-01 01:09:50', 'US_NeonHum': None, 'US_ProvidentEntertainment': None, 'US_ProvidentRCA': '2019-09-20 15:58:52', 'US_ProvidentFilms': '2018-01-12 21:39:59', 'US_ProvidentLabelGroup': '2018-12-01 01:54:18', 'US_RCA': '2018-01-01 00:00:54', 'US_Records': None, 'US_Red': '2019-02-01 13:41:22', 'US_SamePlate': '2018-04-23 18:30:36', 'US_SonyMusicNow': '2019-09-16 18:04:00', 'US_SonyMusicPodcasts': None, 'US_ThreeUncannyFour': None, 'US_PalmTreeRecords': None }, 'theorchard': { 'US': '2019-02-01 13:41:22' } } FIRST_AVAILABLE_DATE_MEMBERSLOGIN = { 'sme': { 'CenturyMediaRecords': '2019-01-27 21:33:05', 'AR': '2019-01-27 17:57:09', 'Asia': '2019-07-30 09:45:21', 'AU_NZ': '2019-01-27 22:49:06', 'AT': '2019-02-18 12:13:24', 'BE': '2019-01-27 17:55:43', 'BR': '2019-01-27 18:01:25', 'BR_FiltrStore': None, 'CA': '2019-01-28 01:56:57', 'CentralAmericaCaribbean': '2019-05-24 20:35:13', 'CL': '2019-01-28 19:50:39', 'CO': '2019-01-27 21:51:34', 'DK': '2019-01-29 15:50:53', 'US': None, 'FI': '2019-10-22 13:29:29', 'FR': '2019-10-15 15:21:57', 'DE': '2019-01-27 16:22:13', 'HK_LiquidState': None, 'IN': '2019-02-28 17:06:40', 'ID': '2019-01-31 10:28:35', 'IE': '2019-03-13 15:54:10', 'IT': '2019-01-27 19:38:34', 'MX': '2019-01-27 16:10:50', 'NL': '2019-01-29 20:25:10', 'NO': '2019-01-27 17:49:37', 'PE': '2019-12-02 23:06:31', 'PL': '2019-02-12 16:04:25', 'PT': '2019-02-01 18:17:46', 'ZA': '2019-02-01 10:07:16', 'ES': '2019-01-27 16:16:58', 'SE': '2019-01-28 11:33:35', 'CH': '2019-01-28 08:13:30', 'TR': '2019-01-30 11:37:00', 'UK': '2019-01-27 15:32:33', 'ClassicalInternational': None, 'WestAfrica_NG': None, 'US_Arista': '2019-02-06 22:05:22', 'US_Columbia': '2019-01-27 16:30:47', 'US_CRM_ProdSupport': '2019-01-30 11:04:17', 'US_Epic': '2019-01-27 16:02:19', 'US_GDB_Sales': '2019-04-03 20:54:58', 'US_LatinIberia': None, 'US_LatinUS': '2019-01-27 15:21:04', 'US_Legacy': '2019-01-27 15:35:20', 'US_Masterworks': '2019-02-21 21:42:40', 'US_MonumentRecords': '2019-01-27 15:35:05', 'US_Nashville': '2019-01-27 15:24:53', 'US_NeonHum': None, 'US_ProvidentEntertainment': None, 'US_ProvidentRCA': '2019-09-20 15:59:45', 'US_ProvidentFilms': '2019-02-05 17:47:43', 'US_ProvidentLabelGroup': '2019-01-27 17:50:53', 'US_RCA': '2019-01-27 16:37:52', 'US_Records': None, 'US_Red': '2019-01-27 17:11:23', 'US_SamePlate': '2019-02-09 15:55:19', 'US_SonyMusicNow': '2019-09-16 19:38:50', 'US_SonyMusicPodcasts': None, 'US_ThreeUncannyFour': None, 'US_PalmTreeRecords': '2019-01-27 16:37:52' }, 'theorchard': { 'US': '2019-01-29 04:18:52' } } FIRST_AVAILABLE_DATE_MEMBERSEXTENDED = { 'sme': { 'CenturyMediaRecords': '2017-10-23 04:00:19', 'AR': '2018-02-13 00:00:08', 'Asia': '2017-12-10 03:00:10', 'AU_NZ': '2017-10-29 10:00:51', 'AT': '2018-10-01 10:00:52', 'AT_MediaSendsNoDOI': None, 'BE': '2017-10-16 18:00:15', 'BR': '2017-10-24 04:00:34', 'BR_FiltrStore': '2020-11-04 16:00:02', 'CA': '2017-10-13 18:00:12', 'CentralAmericaCaribbean': '2019-05-27 18:30:43', 'CL': '2017-11-27 16:02:35', 'CO': '2017-12-11 05:00:12', 'DK': '2018-06-04 22:00:29', 'US': '2019-03-16 23:06:19', 'FI': '2018-05-31 18:02:09', 'FR': '2017-12-16 17:00:14', 'DE': '2017-10-12 15:00:16', 'DE_PromotionNoDOI': None, 'HK_LiquidState': None, 'IN': '2019-02-28 20:24:05', 'ID': '2019-01-31 19:10:05', 'IE': '2018-09-24 22:40:23', 'IT': '2017-10-25 17:00:13', 'MX': '2017-10-06 17:52:08', 'NL': '2017-10-11 08:00:10', 'NO': '2017-11-01 16:00:26', 'PE': '2020-04-29 19:00:23', 'PL': '2017-10-29 20:00:44', 'PT': '2018-11-20 01:10:06', 'ZA': '2018-04-12 01:00:08', 'ES': '2017-10-20 04:00:03', 'SE': '2018-05-17 15:02:09', 'CH': '2017-10-23 04:00:08', 'TR': '2018-02-13 23:00:14', 'UK': '2017-10-10 11:01:05', 'ClassicalInternational': '2021-02-18 07:00:06', 'WestAfrica_NG': '2020-11-16 12:06:01', 'US_Arista': '2019-02-14 03:30:09', 'US_Columbia': '2017-10-11 22:00:16', 'US_CRM_ProdSupport': '2017-10-10 21:00:06', 'US_Epic': '2017-10-10 19:00:15', 'US_GDB_Sales': '2018-08-03 21:00:02', 'US_LatinIberia': '2021-06-03 17:00:02', 'US_LatinUS': '2017-10-14 00:00:18', 'US_Legacy': '2017-10-07 13:00:13', 'US_Masterworks': '2017-12-06 03:00:16', 'US_MonumentRecords': '2017-10-31 03:00:32', 'US_Nashville': '2017-10-07 17:00:10', 'US_NeonHum': '2020-03-24 21:53:10', 'US_ProvidentEntertainment': '2020-09-05 17:31:06', 'US_ProvidentRCA': '2018-06-08 01:00:21', 'US_ProvidentFilms': '2018-01-19 05:01:33', 'US_ProvidentLabelGroup': '2018-10-12 11:01:54', 'US_RCA': '2017-10-06 20:00:05', 'US_Records': '2020-03-27 17:00:07', 'US_Red': '2017-10-30 05:00:21', 'US_SamePlate': '2018-06-19 23:02:12', 'US_SonyMusicNow': '2019-09-20 21:46:14', 'US_SonyMusicPodcasts': '2021-01-21 20:00:04', 'US_ThreeUncannyFour': '2020-10-30 17:59:12', 'US_PalmTreeRecords': None, 'US_ThreeSixZero': None, }, 'theorchard': { 'US': '2017-12-02 01:01:09' } } 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) 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() } already_exsisted_uow = pg.query(slz.UnitOfWork).filter(slz.UnitOfWork.report_date <= start_date.date()).with_entities( slz.UnitOfWork.unit_of_work_code ).all() already_exsisted_uow = {res[0] for res in already_exsisted_uow} 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 ) # skip already created if unit_of_work_code in already_exsisted_uow: logging.info("Already exists, skip: %s", unit_of_work_code) processed.append(True) continue 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.IDLE.value, completeness_status=slz.CompletenessStatusEnum.ACTIVE.value, next_run_at=now, created_at=now, last_updated_at=now, is_force_complete=False, priority=slz.UnitOfWorkPriorityEnum.PRIORITY_6.value, unit_of_work_type=slz.UnitOfWorkTypeEnum.BACKFILL.value, ) 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(): content_status_record = slz.ContentStatus( content_name=f"{context}_{start_date.strftime('%Y%m%d')}.json.gz", context=context, content_status=slz.ContentStatusEnum.ON_HOLD, 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 if uow_record.content_statuses: logging.info("Add: %s", unit_of_work_code) try: pg.add(uow_record) pg.add_all(uow_record.content_statuses) pg.commit() except Exception as err: logging.error("Error during commit: %s", str(err)) pg.rollback() processed.append(True) else: processed.append(False) start_date -= td if not any(processed): break 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)