import logging from datetime import datetime from oa_contract_ingest_to_abacus import oa_contract_json_data_load, oa_contract_json_to_s3_upload, \ oa_contract_json_parse, oa_contract_lambda_execute from oa_contract_ingest_to_abacus.utils import _get_min_contract_id, _get_max_contract_id log = logging.getLogger(__name__) def ingest_contract_to_abacus( account_id: list[int], skip_abacus_import: bool, import_user: str, jira_ticket: str, lambda_function_name: str, s3_bucket: str, ) -> None: if s3_bucket != 'qa-abacus-json-contract-file-import' and s3_bucket != 'uat-abacus-json-contract-file-import': raise Exception('Only QA env is supported at the moment. Safety measures.') # TODO remove when ready to test in Prod log.info("Starting ingestion process...") account_id_to_json_strs = oa_contract_json_data_load.fetch_account_jsons(account_id) log.info(f"Fetched JSON data for account IDs: {account_id}. Total number of items is: {len(account_id_to_json_strs)}") if account_id_to_json_strs: account_id_to_jsons = oa_contract_json_parse.parse_json(account_id_to_json_strs) account_id_to_jsons = sorted(account_id_to_jsons, key=lambda x: int(x[1].get('contract_id', 0))) min_contract_id = _get_min_contract_id(account_id_to_jsons) max_contract_id = _get_max_contract_id(account_id_to_jsons) log.info(f"Computed min_contract_id: {min_contract_id}, max_contract_id: {max_contract_id}") execution_datetime = datetime.now().strftime('%Y%m%d_%H%M%S_%f') batche_keys = oa_contract_json_to_s3_upload.upload_file_in_batches( account_id_to_jsons, f'oa_upload_{execution_datetime}_{min_contract_id}_{max_contract_id}', s3_bucket ) log.info(f'Uploaded batches to S3: {batche_keys}') if not skip_abacus_import: log.info('Starting running lambdas for batches.') oa_contract_lambda_execute.run_lambda_for_batches( batche_keys, import_user, jira_ticket, lambda_function_name, ) else: log.info("No JSON strings found for provided account IDs.")