import re import io from datetime import datetime, date from typing import List, Optional import config from workers.ccp.repository import CCPImportRepository from services.s3_client import S3Client from workers.base_worker import BaseWorker from utils import db_utils, encoding from workers.prs.models import PRSDataFile from models.projects_import_info import ProjectsImportInfo FILE_REGEX = r"^ccp\/ccp_(?Pgrps_projects|po_details_UK-IE|po_details_GSA|expense_hierarchy|budgets|po_details_ibe-it-nav-po-tr|po_details_ca-latam-us)_(?P.*)\.tsv" # noqa: E501 FILE_ORDER = { "grps_projects": 1, "budgets": 2, "expense_hierarchy": 3, "po_details_uk-ie": 4, "po_details_gsa": 5, "po_details_ibe-it-nav-po-tr": 6, "po_details_ca-latam-us": 7 } TABLE_NAME_MAPS = { "grps_projects": "RawCCPProjects", "budgets": "RawCCPBudgets", "expense_hierarchy": "RawCCPExpenseHierarchy", "po_details_uk-ie": "RawCCPPurchaseOrders", "po_details_gsa": "RawCCPPurchaseOrders", # "po_details_ibe-it-nav-po-tr": "RawCCPPurchaseOrders", # "po_details_ca-latam-us": "RawCCPPurchaseOrders", } class CCPDataImportError(Exception): details: str def __init__(self, details: str): self.details = details class CCPDataImporter(BaseWorker): client: S3Client repository = CCPImportRepository() worker_name = "CCPDataImporter" def should_log_exceptions(self): return True def execute(self): self.logger.info("CCP Data import", "Started") self.client = S3Client(config.SFTP_BUCKET) self.__prepare_database_for_import() raw_ccp_purchase_orders_number = self.repository.count_raw_ccp_purchase_orders() self.logger.info("CCP Data import", f"RawCCPPurchaseOrders after preparing db {raw_ccp_purchase_orders_number}") last_import = ProjectsImportInfo.get_imported_ccp_date() or date.min aws_files = self.__get_file_list() most_recent_files = sorted(aws_files, key=lambda x: x.timestamp, reverse=True)[:5] files_to_import = list( filter(lambda x: x.type.lower() in TABLE_NAME_MAPS and x.timestamp.date() > last_import, most_recent_files) ) if not files_to_import: files_names = ", ".join([file.file_name for file in aws_files]) error_details = f"No data files found since {last_import}. Got: {files_names}" self.logger.error("CCP Data import error", error_details) raise CCPDataImportError(error_details) for item in files_to_import: file_object = self.client.get_object(item.file_name, encoding.UTF8).replace("\n\r\n", "\n") with io.StringIO(file_object) as file: if item.is_update_strategy(): db_utils.update_data_from( file=file, table_name=TABLE_NAME_MAPS[item.type.lower()], without_header=True ) else: db_utils.import_data_from( file=file, table_name=TABLE_NAME_MAPS[item.type.lower()], without_header=True ) ProjectsImportInfo.update_imported_ccp_date() self.logger.success("CCP Data import", "Finished!") def __get_file_list(self) -> List[PRSDataFile]: files = self.client.get_objects_list("ccp/") data = [self.__map_from_file(file["Key"]) for file in files] data = [item for item in data if item is not None] data.sort(key=lambda x: FILE_ORDER[x.type.lower()]) return data def __map_from_file(self, file_name: str) -> Optional[PRSDataFile]: match = re.compile(FILE_REGEX, re.IGNORECASE).match(file_name) if not match: if not file_name.endswith(".trg"): self.logger.error("CCP Data file error", f"Invalid file name: {file_name}") return None type = match.group("type") timestamp = match.group("timestamp") try: timestamp = datetime.strptime(timestamp, "%Y%m%d") except ValueError: self.logger.error("CCP Data file error", f"Invalid date format: {timestamp}") return None return PRSDataFile(type, timestamp, file_name) def __prepare_database_for_import(self): self.repository.truncate_raw_ccp_purchase_orders()