import logging from tadas.platform import metrics as metrics_utils from tadas.domain.constants import REPORT_FINAL_TABLES from tadas.monitoring import sensors from tadas.platform import context as contexts from tadas.snowflake import snowflake_publish, tables from tadas.platform import config from tadas.platform import dynamodb as dynamodb_utils from tadas.platform import jenkins as jenkins_utils logger = logging.getLogger(__name__) def publish_to_snowflake(report, model_config): source_table = tables.get_table_name('output', model_config.MODEL_VERSION, report) target_table = REPORT_FINAL_TABLES[report] copy_result = snowflake_publish.copy_from_model_outputs_to_target( source_table=source_table, target_table=target_table, force=config.get('FORCE_PUBLISH'), ) if not copy_result: logger.info(f"No new data was delivered to {target_table}") return logger.info(f"New data was delivered to {target_table}") report_date = snowflake_publish.get_source_report_date(source_table) context = contexts.load_context() context['report_date'] = str(report_date) publish_metrics = metrics_utils.Metrics( category='publish', context=context, ) publish_metrics.add_metric('freshness_threshold_seconds', sensors.freshness_threshold().total_seconds()) publish_metrics.add_metric('published_from', source_table) publish_metrics.add_metric('published_to', target_table) publish_metrics.send() dynamodb_utils.set_overall_status( feed_name=contexts.get_feed_name(report=report), date=str(report_date), overall_status=dynamodb_utils.STATUS_INGESTED, ) _trigger_dbt_if_needed() def _trigger_dbt_if_needed(): if not config.get('TADAS_TRIGGER_JENKINS'): return logger.info('Triggering Jenkins job...') jenkins_utils.trigger_dbt_metrics() logger.info('Triggering Jenkins job completed.')