from datetime import datetime from flask_migrate import Migrate from flask_migrate import MigrateCommand from flask_script import Manager import logging from app import app from app import db from connectors.services.connectors_service import ConnectorsService from projects.services.projects_service import ProjectsService from utils.async_helpers import fire_and_forget from utils.reporting.worker_logger import WorkerLogger from services.artist_team_automation.artist_team_automation_service import ArtistTeamAutomationService from workers.prs.prs_project_budget_recalculation_worker import PRSProjectBudgetsRecalculationWorker from workers.prs.prs_budget_worker import PRSBudgetWorker from workers.prs.prs_purchase_orders_worker import PRSPurchaseOrdersWorker from workers.prs.prs_projects_worker import PRSProjectsWorker from workers.prs.prs_data_worker import PRSDataImporter from workers.prs.prs_elastic_indexer_worker import PRSElasticIndexerWorker from workers.prs.stream_rates.prs_stream_rates_data_worker import PRSStreamRatesDataImporter from workers.prs.stream_rates.prs_streams_rates_worker import PRSStreamRatesWorker from workers.prs.prs_budget_group_category_worker import PRSBudgetGroupCategoryWorker from workers.product_families_worker import ProductFamiliesWorker from workers.confidential_projects_worker import ConfidentialProjectWorker from workers.unknown_artists_worker import UnknownArtistsWorker from workers.artist_team_automation_worker import ArtistTeamAutomationWorker from workers.labels_copy_worker import LabelsCopyWorker from workers.ccp.ccp_data_worker import CCPDataImporter from workers.ccp.ccp_projects_worker import CCPProjectsWorker from workers.ccp.ccp_budget_group_category_worker import CCPBudgetGroupCategoryWorker from workers.ccp.ccp_purchase_orders_worker import CCPPurchaseOrdersWorker from workers.ccp.ccp_budgets_worker import CCPBudgetWorker from workers.ccp.ccp_budgets_recalculation_worker import CCPBudgetsRecalculationWorker from worker import RQWorker, RQWorkerCleaner from rq import Retry migrate = Migrate(app, db) manager = Manager(app) manager.add_command("db", MigrateCommand) manager.add_command("start_queue_worker", RQWorker) manager.add_command("clean_worker_queue", RQWorkerCleaner) @manager.command def reindex_elastic(model=None): def model_in_models_list(name: str, models_list: list): for m in models_list: if m.__name__ == name: return m return None from utils.elastic_search.indexer import ElasticSearchIndexer elastic_indexer = ElasticSearchIndexer() if model is not None: index_model = model_in_models_list(model, elastic_indexer.models_to_index) if index_model is not None: elastic_indexer.reindex_database(index_model) else: raise ValueError(f"{model} not valid model name, should be one of {[m.__name__ for m in elastic_indexer.models_to_index]}") else: elastic_indexer.reindex_database() @manager.command def import_gras_data(): from utils.snowflake.gras_data_importer import GRASDataImporter if app.debug: logging.basicConfig(level=logging.INFO) gras_data_importer = GRASDataImporter() gras_data_importer.import_projects() @manager.command def import_campaigns_data(): from utils.snowflake.google_data.importer import GoogleCampaignsAdsImporter from utils.snowflake.facebook_data.importer import FacebookCampaignsImporter from utils.snowflake.tiktok_data.importer import TikTokCampaignsImporter if app.debug: logging.basicConfig(level=logging.INFO) google_importer = GoogleCampaignsAdsImporter() google_importer.import_campaigns() fb_importer = FacebookCampaignsImporter() fb_importer.import_campaigns() tiktok_importer = TikTokCampaignsImporter() tiktok_importer.import_campaigns() @manager.command def __import_google_campaigns_data(): from utils.snowflake.google_data.importer import GoogleCampaignsAdsImporter if app.debug: logging.basicConfig(level=logging.INFO) google_importer = GoogleCampaignsAdsImporter() google_importer.import_campaigns() @manager.command def __import_facebook_campaigns_data(): from utils.snowflake.facebook_data.importer import FacebookCampaignsImporter if app.debug: logging.basicConfig(level=logging.INFO) importer = FacebookCampaignsImporter() importer.import_campaigns() @manager.command def __import_tiktok_campaigns_data(): from utils.snowflake.tiktok_data.importer import TikTokCampaignsImporter if app.debug: logging.basicConfig(level=logging.INFO) importer = TikTokCampaignsImporter() importer.import_campaigns() @manager.command def import_ad_sets_data(): from utils.snowflake.google_data.importer import GoogleAdSetsImporter from utils.snowflake.facebook_data.importer import FacebookAdSetsImporter from utils.snowflake.tiktok_data.importer import TikTokAdSetsImporter if app.debug: logging.basicConfig(level=logging.INFO) google_importer = GoogleAdSetsImporter() google_importer.import_ads() fb_importer = FacebookAdSetsImporter() fb_importer.import_ads() tiktok_importer = TikTokAdSetsImporter() tiktok_importer.import_ads() @manager.command def __import_google_ad_sets_data(): from utils.snowflake.google_data.importer import GoogleAdSetsImporter if app.debug: logging.basicConfig(level=logging.INFO) importer = GoogleAdSetsImporter() importer.import_ads() @manager.command def __import_facebook_ad_sets_data(): from utils.snowflake.facebook_data.importer import FacebookAdSetsImporter if app.debug: logging.basicConfig(level=logging.INFO) importer = FacebookAdSetsImporter() importer.import_ads() @manager.command def __import_tiktok_ad_sets_data(): from utils.snowflake.tiktok_data.importer import TikTokAdSetsImporter if app.debug: logging.basicConfig(level=logging.INFO) importer = TikTokAdSetsImporter() importer.import_ads() @manager.command def import_artists_moments(): from utils.snowflake.gras_data.batch_moments_importer import BatchMomentsImporter if app.debug: logging.basicConfig(level=logging.INFO) importer = BatchMomentsImporter() importer.import_moments() @manager.command def import_prs_data(): data_worker = PRSDataImporter() data_worker_job = data_worker.perform_async() process_prs_data(data_worker_job) @manager.command def import_ccp_data(): data_worker = CCPDataImporter() data_worker_job = data_worker.perform_async() process_ccp_data(data_worker_job) @manager.command def import_projects_data(): import_prs_worker = PRSDataImporter() import_ccp_worker = CCPDataImporter() prs_worker_job = import_prs_worker.perform_async(retry=Retry(max=3, interval=[10, 30, 60])) final_process_prs_job = process_prs_data(prs_worker_job) ccp_worker_job = import_ccp_worker.perform_async( depends_on=final_process_prs_job, retry=Retry(max=3, interval=[10, 30, 60]), ) final_process_ccp_job = process_ccp_data(ccp_worker_job) product_families_job = ProductFamiliesWorker().perform_async(depends_on=final_process_ccp_job) unknown_artists_worker_job = UnknownArtistsWorker().perform_async(depends_on=product_families_job) elastic_indexer_job = PRSElasticIndexerWorker().perform_async(depends_on=unknown_artists_worker_job) ArtistTeamAutomationWorker().perform_async(depends_on=elastic_indexer_job) @manager.command def process_ccp_data(data_worker_job=None): ccp_projects_worker = CCPProjectsWorker() if data_worker_job: ccp_projects_worker_job = ccp_projects_worker.perform_async(depends_on=data_worker_job) else: ccp_projects_worker_job = ccp_projects_worker.perform_async() CCPBudgetGroupCategoryWorker().perform_async(depends_on=ccp_projects_worker_job) confidential_projects_worker_job = ConfidentialProjectWorker().perform_async(depends_on=ccp_projects_worker_job) ccp_purchase_orders_worker_job = CCPPurchaseOrdersWorker().perform_async( depends_on=confidential_projects_worker_job ) ccp_budgets_worker_job = CCPBudgetWorker().perform_async(depends_on=ccp_purchase_orders_worker_job) return CCPBudgetsRecalculationWorker().perform_async(depends_on=ccp_budgets_worker_job) @manager.command def refresh_delphi_triggers(): from migrations.fixtures.delphi_tables.v1.data_refresher import DelphiTablesDataRefresher DelphiTablesDataRefresher().refresh(db.session) @manager.command def process_prs_data(data_worker_job=None): projects_worker = PRSProjectsWorker() if data_worker_job: projects_worker_job = projects_worker.perform_async(depends_on=data_worker_job) else: projects_worker_job = projects_worker.perform_async() PRSBudgetGroupCategoryWorker().perform_async(depends_on=projects_worker_job) confidential_projects_worker_job = ConfidentialProjectWorker().perform_async(depends_on=projects_worker_job) prs_po_worker_job = PRSPurchaseOrdersWorker().perform_async(depends_on=confidential_projects_worker_job) prs_budgets_worker = PRSBudgetWorker().perform_async(depends_on=prs_po_worker_job) return PRSProjectBudgetsRecalculationWorker().perform_async(depends_on=prs_budgets_worker) @manager.command def import_stream_rates_data(): stream_rates_data_worker = PRSStreamRatesDataImporter() stream_rates_data_worker_job = stream_rates_data_worker.perform_async() PRSStreamRatesWorker().perform_async(depends_on=stream_rates_data_worker_job) @manager.command def __import_product_families(): product_families_worker = ProductFamiliesWorker() product_families_worker.perform_async() @manager.command def __gras_protected_projects(): ConfidentialProjectWorker().perform_async() @manager.command def __import_unknown_artists(): UnknownArtistsWorker().perform_async() @manager.command def autoclaim_available_projects(): ArtistTeamAutomationService().autoclaim_all_available_projects() @manager.command def check_projects_import(): logger = WorkerLogger() projects_service = ProjectsService() logger.info("Projects import check started", f"At {datetime.now()}") try: projects_service.check_projects_import() except Exception as error: logger.error("Projects import check failed", f"{error}") return logger.success("Project import check is successful", f"At {datetime.now()}") @manager.command def copy_labels_data(): LabelsCopyWorker().perform_async() @manager.command def register_fivetran_webhook(): logger = WorkerLogger() connectors_service = ConnectorsService() logger.info("Creating Fivetran webhook", f"At {datetime.now()}") fire_and_forget(connectors_service.register_webhook()) logger.success("Creating Fivetran webhoook is finished check manually status", f"At {datetime.now()}") # check webhook success by manually Fivetran API request if __name__ == "__main__": manager.run()