import traceback from typing import List from campaigns.services.campaigns_history_service import CampaignHistoryService from db import db from projects.repositories.projects_repository import ProjectsRepository from utils.list_utils import chunks, range_chunks from models.territories import Territory from workers.base_worker import BaseWorker from workers.prs.mappers.moment_mapper import MomentMapper from workers.prs.mappers.purchase_order_mapper import PurchaseOrderMapper from workers.prs.repository import PRSImportRepository from models.raw_prs_purchase_orders_map import PRSMappingEntityType from models.prs_purchase_order import PRSPurchaseOrderStatus, PRSPurchaseOrder from campaigns.repositories.campaigns_repository import CampaignsRepository from services.territory.territories_repository import TerritoriesRepository UNPROCESSABLE_STATUS = [ PRSPurchaseOrderStatus.REJECTED.value, PRSPurchaseOrderStatus.REMOVED.value, PRSPurchaseOrderStatus.FAILED.value, ] DEFAULT_GENDERS = [1, 2, 3] DEFAULT_TERRITORIES = [0] def is_none(value): return value is None or isinstance(value, type(None)) class PRSPurchaseOrdersWorker(BaseWorker): repository = PRSImportRepository() campaigns_repository = CampaignsRepository() territories_repository = TerritoriesRepository() projects_repository = ProjectsRepository() history_service: CampaignHistoryService = CampaignHistoryService() worker_name = "PRSPurchaseOrdersWorker" default_territories: List[Territory] def should_log_exceptions(self): return True def execute(self, *args): new_providers = self.repository.get_new_providers() for chunk in chunks(new_providers, 500): db.session.bulk_save_objects(list(map(lambda x: self.repository.create_provider_with_name(x), chunk))) db.session.commit() raw_purchase_orders_count = self.repository.get_raw_prs_purchase_order_count() self.default_territories = self.territories_repository.get_territories_by_ids(DEFAULT_TERRITORIES) self.logger.info("PRS POs import", f"Importing {raw_purchase_orders_count} POs...") po_per_iteration = 50000 try: for chunk in range_chunks(range(0, raw_purchase_orders_count), po_per_iteration): raw_purchase_orders = self.repository.get_raw_prs_purchase_order(po_per_iteration, chunk.start) self.__execute_po_import(raw_purchase_orders) except Exception as error: self.logger.error("PRS POs import error", f"{error}\n{traceback.format_exc()}") return self.logger.success("PRS POs import", "Finished!") def __execute_po_import(self, purchase_orders): projects_ids = set() for result in purchase_orders: project_id = result.project_id po = self.__map_purchase_order(project_id, result) if po.id is None or db.session.is_modified(po): projects_ids.add(project_id) db.session.add(po) db.session.flush() mapping_type = result.mapping_entity_type if not mapping_type: continue if mapping_type == PRSMappingEntityType.MOMENT.value: self.__map_moment_from_po(po, result.moment_type_id, result.artist_id, result.artist_external_id) self.projects_repository.touch_projects_last_edit(list(projects_ids)) db.session.commit() def __map_moment_from_po(self, purchase_order, type_id, artist_id: int, artist_external_id: str): moment = self.__map_moment(purchase_order, type_id, self.default_territories, artist_id, artist_external_id) db.session.add(moment) def __map_purchase_order(self, project_id, purchase_order_model) -> PRSPurchaseOrder: return PurchaseOrderMapper(project_id=project_id, purchase_order=purchase_order_model).build() def __map_moment( self, purchase_order_model: PRSPurchaseOrder, moment_type_id: int, territories: List[Territory], artist_id: int, artist_external_id: str ): return MomentMapper( purchase_order=purchase_order_model, moment_type_id=moment_type_id, artist_id=artist_id, artist_external_id=artist_external_id, territories=territories ).build()